diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 51a950348..a69573aa7 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -80,6 +80,8 @@ jobs: Path("scripts/cmux_fleet_bootstrap.py"), Path("scripts/test-cmux-fleet.py"), Path("scripts/test-cmux-fleet-bootstrap.py"), + Path("scripts/cmux_workload_request.py"), + Path("scripts/test-cmux-workload-request.py"), Path("scripts/test-verify-focused.py"), Path("scripts/verify-changed-rustfmt"), Path("scripts/test-hot-run.py"), @@ -121,6 +123,7 @@ jobs: python3 scripts/test-owned-linux-admission.py python3 scripts/test-cmux-fleet.py python3 scripts/test-cmux-fleet-bootstrap.py + python3 scripts/test-cmux-workload-request.py python3 scripts/test-verify.py python3 scripts/test-verify-changed-rustfmt.py python3 scripts/test-front-door.py diff --git a/docs/CMUX_WORKLOAD_PROFILES.md b/docs/CMUX_WORKLOAD_PROFILES.md new file mode 100644 index 000000000..8444d8dde --- /dev/null +++ b/docs/CMUX_WORKLOAD_PROFILES.md @@ -0,0 +1,114 @@ +# CMUX repository-owned workload profiles + +CMUX owns the meaning of CMUX workload profiles. Glaeda may admit, place, execute, cache, time, and recover those workloads without copying their command lines or redefining pass/fail semantics. + +The paired CMUX implementation is tracked by `manaflow-ai/cmux#13411`. + +## Request boundary + +`scripts/cmux_workload_request.py` accepts a bounded semantic request containing: + +- exact `manaflow-ai/cmux` commit and tree; +- CMUX profile ID and explicit semantic generation; +- benchmark state class; +- bounded integer semantic parameters, such as an app-host shard number; +- caller correlation and an advisory reuse hint. + +It rejects caller-controlled argv, shell text, cwd, host paths, backend/machine choice, resource values, cache roots, attempt identity, or cleanup authority. + +A planned request resolves only this fixed adapter: + +```json +{ + "kind": "cmux-repository-profile/v1", + "repository_runner": "scripts/ci/cmux_workload_profile.py", + "semantic_result_contract": "cmux-workload-result/v1" +} +``` + +Glaeda does not keep a copy of the CMUX registry. Profile validity is checked after exact source materialization by the runner inside that CMUX tree. + +## Physical execution + +A physical adapter follows this sequence: + +1. materialize the exact requested CMUX commit/tree; +2. choose an eligible backend and admitted cache/state class from Glaeda-owned observations; +3. run the checked-in CMUX profile runner's `plan` command and require the same source, profile ID/generation, state class, and semantic parameters; +4. allocate a task-private state root and any reviewed runtime inputs; +5. run the same CMUX profile runner and retain the exact canonical `cmux-workload-result/v1`; +6. wrap its digest and terminal state in Glaeda's physical receipt. + +The generic repository adapter may translate the semantic request into the fixed CMUX runner interface. It does not reconstruct Xcode, test-selector, package, release, or developer-build commands. + +Runtime-input paths stay physical. Their content identities are published by the CMUX semantic result when the profile contract requires them. + +## Identity + +Three identities remain separate: + +- caller correlation: external request/work references; +- semantic execution binding: exact source + CMUX profile/generation + state class + semantic parameters + fixed adapter generation; +- physical attempt: Glaeda backend, node, leases, cache placement, resources, timing, cleanup, and recovery state. + +Caller correlation and reuse hints do not mint another execution binding. A profile generation change does. + +## Result correlation + +`observe` accepts only canonical `cmux-workload-result/v1` bytes and verifies: + +- exact repository/commit/tree; +- exact profile ID/generation; +- exact semantic parameters; +- requested benchmark state class; +- presence of the CMUX semantic validator, artifact collection, and cleanup evidence. + +It validates the shared `cmux-workload-result/v1` integrity envelope before projection: the result schema is closed, nested identity/evidence shapes are typed and bounded by the document ceiling, the toolchain identity must match its observations, timestamps and cleanup are self-consistent, and a `passed` result must have exit zero, complete required-artifact validation, and clean process settlement. These are interface-integrity checks; Glaeda does not re-run or reinterpret CMUX's profile-specific validators. + +It then projects the CMUX terminal vocabulary without redefining profile semantics: + +- `passed` -> `succeeded`; +- `failed` -> `failed`; +- `timed_out` -> `timed_out`; +- `ambiguous` -> `ambiguous`. + +The outer observation stores the exact CMUX result digest. Glaeda-specific machine and execution evidence belongs beside that digest in the physical receipt. + +## Fleet roles and routing + +The first CMUX mapping is: + +```yaml +cmux_linux_ci: + acceptance: + profile: cmux.ci.guard + generation: 1 + +cmux_macos_native_build: + acceptance: + profile: cmux.macos.dev-check + generation: 1 +``` + +Role acceptance consumes the current profile identity from CMUX rather than a Glaeda-owned recipe. Routing experiments can compare backends only when the CMUX semantic comparison identity and benchmark context agree. + +This gives #148 a repository-owned named-profile source, #546 comparable routing evidence, #547 a stable semantic operation for optimization, and #1056 an exact acceptance workload. #1057's caller/orchestrator identity remains independent from this workload identity. + +## Synthetic proof + +`docs/experiments/cmux-workload-profile/` contains a request and its deterministic planned result. It proves the GitHub/repository-only request boundary without claiming physical execution. + +Run: + +```sh +python3 scripts/cmux_workload_request.py plan \ + < docs/experiments/cmux-workload-profile/request.json \ + > /tmp/cmux-workload-plan.json + +cmp /tmp/cmux-workload-plan.json \ + docs/experiments/cmux-workload-profile/plan.json + +python3 scripts/test-cmux-workload-request.py +``` + +Physical Mac acceptance remains a separate approved fleet experiment. diff --git a/docs/experiments/cmux-workload-profile/README.md b/docs/experiments/cmux-workload-profile/README.md new file mode 100644 index 000000000..d61aa908a --- /dev/null +++ b/docs/experiments/cmux-workload-profile/README.md @@ -0,0 +1,17 @@ +# CMUX workload profile synthetic plan + +This fixture proves Glaeda can accept a caller-neutral CMUX semantic workload request while leaving CMUX semantics in the CMUX repository. + +`request.json` freezes synthetic exact source identity plus `cmux.ci.guard@1` and the `cold` state class. `plan.json` resolves only the fixed `cmux-repository-profile/v1` adapter and grants no execution, placement, resource-override, or redispatch authority. + +The OIDs are deliberately synthetic; this fixture exercises the request contract without claiming that a CMUX checkout or physical worker ran. + +Regenerate and compare: + +```sh +python3 scripts/cmux_workload_request.py plan \ + < docs/experiments/cmux-workload-profile/request.json \ + > /tmp/cmux-workload-plan.json +cmp /tmp/cmux-workload-plan.json \ + docs/experiments/cmux-workload-profile/plan.json +``` diff --git a/docs/experiments/cmux-workload-profile/plan.json b/docs/experiments/cmux-workload-profile/plan.json new file mode 100644 index 000000000..a300edde6 --- /dev/null +++ b/docs/experiments/cmux-workload-profile/plan.json @@ -0,0 +1 @@ +{"authority":{"authorizes_execution":false,"authorizes_host_selection":false,"authorizes_redispatch":false,"authorizes_resource_override":false},"benchmark_state_class":"cold","correlation":{"work_ref":"cmux:work:fixture-1"},"document_type":"glaeda-cmux-workload-plan","execution_binding_sha256":"sha256:45ac07054ec2ba65cb46d039c879b6b973c33e2fe3f2290033f09981bfe1b8ff","external_request_ref":"cmux:workload:fixture-1","parameters":{},"physical_preflight":["materialize_exact_source","cmux_runner_plan_matches_frozen_source_and_profile","cmux_runner_owns_semantic_validation"],"profile":{"generation":1,"id":"cmux.ci.guard"},"request_sha256":"sha256:2ff1cc2403bfd2226d0e610fca06d74c0de221f18a252452b6cf58cc742a5d81","resolved_adapter":{"kind":"cmux-repository-profile/v1","repository_runner":"scripts/ci/cmux_workload_profile.py","semantic_result_contract":"cmux-workload-result/v1"},"schema_version":1,"source":{"commit":"1111111111111111111111111111111111111111","repository":"manaflow-ai/cmux","tree":"2222222222222222222222222222222222222222"},"state":"planned"} diff --git a/docs/experiments/cmux-workload-profile/request.json b/docs/experiments/cmux-workload-profile/request.json new file mode 100644 index 000000000..e13bd8a05 --- /dev/null +++ b/docs/experiments/cmux-workload-profile/request.json @@ -0,0 +1 @@ +{"benchmark_state_class":"cold","correlation":{"work_ref":"cmux:work:fixture-1"},"document_type":"glaeda-cmux-workload-request","external_request_ref":"cmux:workload:fixture-1","profile":{"generation":1,"id":"cmux.ci.guard"},"reuse_hint":"no_preference","schema_version":1,"source":{"commit":"1111111111111111111111111111111111111111","repository":"manaflow-ai/cmux","tree":"2222222222222222222222222222222222222222"}} diff --git a/scripts/cmux_workload_request.py b/scripts/cmux_workload_request.py new file mode 100644 index 000000000..776c2ddfc --- /dev/null +++ b/scripts/cmux_workload_request.py @@ -0,0 +1,666 @@ +#!/usr/bin/env python3 +"""Caller-neutral Glaeda adapter for CMUX-owned workload profiles.""" +from __future__ import annotations + +import argparse +import hashlib +import json +import re +import sys +from dataclasses import dataclass +from pathlib import Path +from typing import NoReturn + +REQUEST_TYPE = "glaeda-cmux-workload-request" +PLAN_TYPE = "glaeda-cmux-workload-plan" +OBSERVATION_TYPE = "glaeda-cmux-workload-observation" +SCHEMA_VERSION = 1 +MAX_DOCUMENT_BYTES = 16 * 1024 +CMUX_REPOSITORY = "manaflow-ai/cmux" +CMUX_RUNNER = "scripts/ci/cmux_workload_profile.py" +CMUX_RESULT_CONTRACT = "cmux-workload-result/v1" +OID_RE = re.compile(r"^[0-9a-f]{40}$") +DIGEST_RE = re.compile(r"^sha256:[0-9a-f]{64}$") +PROFILE_RE = re.compile(r"^cmux\.[a-z0-9][a-z0-9.-]{1,80}$") +TOKEN_RE = re.compile(r"^[A-Za-z0-9][A-Za-z0-9._:/@+-]{0,127}$") +PARAM_RE = re.compile(r"^[a-z][a-z0-9_]{0,31}$") +STATE_CLASSES = { + "cold", + "dependency-warm", + "compiler-warm", + "exact-product-reuse", + "resident-hot", +} +SEMANTIC_RESULTS = {"passed", "failed", "timed_out", "ambiguous"} +OUTER_STATES = { + "passed": "succeeded", + "failed": "failed", + "timed_out": "timed_out", + "ambiguous": "ambiguous", +} +AUTHORITY = { + "authorizes_execution": False, + "authorizes_host_selection": False, + "authorizes_resource_override": False, + "authorizes_redispatch": False, +} + + +class ContractRefusal(ValueError): + def __init__(self, code: str, message: str): + super().__init__(message) + self.code = code + + +@dataclass(frozen=True) +class Request: + external_request_ref: str + repository: str + commit: str + tree: str + profile_id: str + profile_generation: int + benchmark_state_class: str + parameters: dict[str, int] + reuse_hint: str | None + work_ref: str | None + + +def canonical_bytes(value: object) -> bytes: + return json.dumps( + value, sort_keys=True, separators=(",", ":"), allow_nan=False + ).encode("utf-8") + + +def sha256(value: bytes) -> str: + return "sha256:" + hashlib.sha256(value).hexdigest() + + +def reject_json_constant(value: str) -> NoReturn: + raise ValueError(value) + + +def token(value: object, label: str) -> str: + if not isinstance(value, str) or TOKEN_RE.fullmatch(value) is None: + raise ContractRefusal("invalid_request", f"{label} is invalid") + return value + + +def decode_json(raw: bytes, ceiling: int = MAX_DOCUMENT_BYTES) -> object: + if len(raw) > ceiling: + raise ContractRefusal("invalid_document", "document exceeds its fixed ceiling") + try: + return json.loads(raw, parse_constant=reject_json_constant) + except (UnicodeError, ValueError) as error: + raise ContractRefusal("invalid_document", "document is not valid JSON") from error + + +def decode_request(raw: bytes) -> Request: + value = decode_json(raw) + required = { + "document_type", + "schema_version", + "external_request_ref", + "source", + "profile", + "benchmark_state_class", + } + optional = {"parameters", "reuse_hint", "correlation"} + if not isinstance(value, dict) or not required <= set(value) <= required | optional: + raise ContractRefusal("invalid_request", "request fields are invalid") + if ( + value["document_type"] != REQUEST_TYPE + or type(value["schema_version"]) is not int + or value["schema_version"] != SCHEMA_VERSION + ): + raise ContractRefusal("unsupported_schema", "request schema is unsupported") + source = value["source"] + if not isinstance(source, dict) or set(source) != {"repository", "commit", "tree"}: + raise ContractRefusal("invalid_request", "source identity is invalid") + if source["repository"] != CMUX_REPOSITORY: + raise ContractRefusal("unsupported_repository", "repository is unsupported") + if ( + not isinstance(source["commit"], str) + or OID_RE.fullmatch(source["commit"]) is None + or not isinstance(source["tree"], str) + or OID_RE.fullmatch(source["tree"]) is None + ): + raise ContractRefusal("invalid_request", "source commit/tree identity is invalid") + profile = value["profile"] + if not isinstance(profile, dict) or set(profile) != {"id", "generation"}: + raise ContractRefusal("invalid_request", "profile identity is invalid") + if not isinstance(profile["id"], str) or PROFILE_RE.fullmatch(profile["id"]) is None: + raise ContractRefusal("invalid_request", "profile id is invalid") + if ( + type(profile["generation"]) is not int + or profile["generation"] < 1 + or profile["generation"] > 2**31 - 1 + ): + raise ContractRefusal("invalid_request", "profile generation is invalid") + state_class = value["benchmark_state_class"] + if not isinstance(state_class, str) or state_class not in STATE_CLASSES: + raise ContractRefusal("invalid_request", "benchmark state class is invalid") + parameters_value = value.get("parameters", {}) + if not isinstance(parameters_value, dict) or len(parameters_value) > 8: + raise ContractRefusal("invalid_request", "profile parameters are invalid") + parameters: dict[str, int] = {} + for name, parameter in parameters_value.items(): + if ( + not isinstance(name, str) + or PARAM_RE.fullmatch(name) is None + or type(parameter) is not int + or parameter < -(2**31) + or parameter > 2**31 - 1 + ): + raise ContractRefusal("invalid_request", "profile parameter is invalid") + parameters[name] = parameter + reuse_hint = value.get("reuse_hint") + if reuse_hint is not None and reuse_hint not in {"no_preference", "prefer_valid_reuse"}: + raise ContractRefusal("invalid_request", "reuse hint is invalid") + work_ref = None + correlation = value.get("correlation") + if correlation is not None: + if not isinstance(correlation, dict) or set(correlation) != {"work_ref"}: + raise ContractRefusal("invalid_request", "correlation is invalid") + work_ref = token(correlation["work_ref"], "work reference") + return Request( + external_request_ref=token( + value["external_request_ref"], "external request reference" + ), + repository=source["repository"], + commit=source["commit"], + tree=source["tree"], + profile_id=profile["id"], + profile_generation=profile["generation"], + benchmark_state_class=state_class, + parameters=parameters, + reuse_hint=reuse_hint, + work_ref=work_ref, + ) + + +def request_document(request: Request) -> dict[str, object]: + document: dict[str, object] = { + "document_type": REQUEST_TYPE, + "schema_version": SCHEMA_VERSION, + "external_request_ref": request.external_request_ref, + "source": { + "repository": request.repository, + "commit": request.commit, + "tree": request.tree, + }, + "profile": { + "id": request.profile_id, + "generation": request.profile_generation, + }, + "benchmark_state_class": request.benchmark_state_class, + } + if request.parameters: + document["parameters"] = request.parameters + if request.reuse_hint is not None: + document["reuse_hint"] = request.reuse_hint + if request.work_ref is not None: + document["correlation"] = {"work_ref": request.work_ref} + return document + + +def execution_binding(request: Request) -> str: + return sha256( + canonical_bytes( + { + "domain": "glaeda-cmux-repository-profile-binding-v1", + "source": { + "repository": request.repository, + "commit": request.commit, + "tree": request.tree, + }, + "profile": { + "id": request.profile_id, + "generation": request.profile_generation, + }, + "benchmark_state_class": request.benchmark_state_class, + "parameters": request.parameters, + "adapter": { + "kind": "cmux-repository-profile/v1", + "runner": CMUX_RUNNER, + "result_contract": CMUX_RESULT_CONTRACT, + }, + } + ) + ) + + +def plan(request: Request) -> dict[str, object]: + result: dict[str, object] = { + "document_type": PLAN_TYPE, + "schema_version": SCHEMA_VERSION, + "external_request_ref": request.external_request_ref, + "request_sha256": sha256(canonical_bytes(request_document(request))), + "execution_binding_sha256": execution_binding(request), + "source": { + "repository": request.repository, + "commit": request.commit, + "tree": request.tree, + }, + "profile": { + "id": request.profile_id, + "generation": request.profile_generation, + }, + "benchmark_state_class": request.benchmark_state_class, + "parameters": request.parameters, + "state": "planned", + "resolved_adapter": { + "kind": "cmux-repository-profile/v1", + "repository_runner": CMUX_RUNNER, + "semantic_result_contract": CMUX_RESULT_CONTRACT, + }, + "physical_preflight": [ + "materialize_exact_source", + "cmux_runner_plan_matches_frozen_source_and_profile", + "cmux_runner_owns_semantic_validation", + ], + "authority": AUTHORITY, + } + if request.work_ref is not None: + result["correlation"] = {"work_ref": request.work_ref} + return result + + +def _exact_keys(value: object, expected: set[str], label: str) -> dict[str, object]: + if not isinstance(value, dict) or set(value) != expected: + raise ContractRefusal("invalid_result", f"{label} fields are invalid") + return value + + +def _cmux_semantic_key(result: dict[str, object]) -> str: + source = result["source"] + profile = result["profile"] + runtime_inputs = result["runtime_input_identities"] + assert isinstance(source, dict) + assert isinstance(profile, dict) + assert isinstance(runtime_inputs, list) + return sha256( + canonical_bytes( + { + "source": { + "repository": source["repository"], + "tree": source["tree"], + }, + "profile": { + "id": profile["id"], + "generation": profile["generation"], + }, + "semantic_validator": result["semantic_validator"], + "parameters": result["parameters"], + "runtime_inputs": [ + { + "name": item["name"], + "class": item["class"], + "identity": item["identity"], + "sha256": item["sha256"], + } + for item in runtime_inputs + if isinstance(item, dict) + ], + } + ) + ) + + +def _cmux_context_key(semantic_key: str, state_class: str, toolchain: str) -> str: + return sha256( + canonical_bytes( + { + "semantic_key": semantic_key, + "state_class": state_class, + "toolchain_identity": toolchain, + } + ) + ) + + +def validate_cmux_result(value: dict[str, object]) -> dict[str, object]: + _exact_keys( + value, + { + "document_type", + "schema_version", + "source", + "profile", + "semantic_validator", + "expected_result_class", + "result", + "parameters", + "runtime_input_identities", + "artifact_identities", + "validation", + "stage_timings", + "resource_summary", + "toolchain", + "benchmark", + "network_class", + "timeout_class", + "cleanup", + "exit_code", + "started_at_unix_millis", + "ended_at_unix_millis", + }, + "CMUX semantic result", + ) + source = _exact_keys( + value["source"], + {"repository", "commit", "tree"}, + "CMUX semantic result source", + ) + profile = _exact_keys( + value["profile"], + {"id", "generation"}, + "CMUX semantic result profile", + ) + validation = _exact_keys( + value["validation"], + {"missing_required_artifact_classes"}, + "CMUX semantic result validation", + ) + benchmark = _exact_keys( + value["benchmark"], + {"state_class", "semantic_comparison_key", "comparison_context_key"}, + "CMUX semantic result benchmark", + ) + cleanup = _exact_keys( + value["cleanup"], + {"state", "process_group_settled"}, + "CMUX semantic result cleanup", + ) + toolchain = _exact_keys( + value["toolchain"], + {"identity", "observations"}, + "CMUX semantic result toolchain", + ) + resource = _exact_keys( + value["resource_summary"], + {"resource_class", "cpu_count", "memory_bytes", "architecture"}, + "CMUX semantic result resources", + ) + if ( + source.get("repository") != CMUX_REPOSITORY + or not isinstance(source.get("commit"), str) + or OID_RE.fullmatch(source["commit"]) is None + or not isinstance(source.get("tree"), str) + or OID_RE.fullmatch(source["tree"]) is None + or not isinstance(profile.get("id"), str) + or PROFILE_RE.fullmatch(profile["id"]) is None + or type(profile.get("generation")) is not int + or profile["generation"] < 1 + or profile["generation"] > 2**31 - 1 + ): + raise ContractRefusal("invalid_result", "CMUX semantic source/profile is invalid") + if ( + value.get("result") not in SEMANTIC_RESULTS + or not isinstance(value.get("semantic_validator"), str) + or not value["semantic_validator"] + or not isinstance(value.get("expected_result_class"), str) + or not value["expected_result_class"] + or not isinstance(value.get("network_class"), str) + or not value["network_class"] + or not isinstance(value.get("timeout_class"), str) + or not value["timeout_class"] + or type(value.get("exit_code")) is not int + or type(value.get("started_at_unix_millis")) is not int + or type(value.get("ended_at_unix_millis")) is not int + or not 0 <= value["started_at_unix_millis"] <= value["ended_at_unix_millis"] < 253402300800000 + ): + raise ContractRefusal("invalid_result", "CMUX semantic result state is invalid") + parameters = value["parameters"] + if ( + not isinstance(parameters, dict) + or len(parameters) > 8 + or any( + not isinstance(name, str) + or PARAM_RE.fullmatch(name) is None + or type(parameter) is not int + for name, parameter in parameters.items() + ) + ): + raise ContractRefusal("invalid_result", "CMUX semantic parameters are invalid") + + runtime_inputs = value["runtime_input_identities"] + if not isinstance(runtime_inputs, list): + raise ContractRefusal("invalid_result", "CMUX runtime inputs are invalid") + seen_inputs: set[str] = set() + for item in runtime_inputs: + entry = _exact_keys( + item, + {"name", "class", "identity", "sha256", "bytes"}, + "CMUX runtime input", + ) + if ( + not isinstance(entry.get("name"), str) + or not entry["name"] + or entry["name"] in seen_inputs + or not isinstance(entry.get("class"), str) + or not entry["class"] + or entry.get("identity") not in {"file-sha256", "parent-tree-sha256"} + or not isinstance(entry.get("sha256"), str) + or DIGEST_RE.fullmatch(entry["sha256"]) is None + or type(entry.get("bytes")) is not int + or entry["bytes"] < 0 + ): + raise ContractRefusal("invalid_result", "CMUX runtime input is invalid") + seen_inputs.add(entry["name"]) + + artifacts = value["artifact_identities"] + if not isinstance(artifacts, list): + raise ContractRefusal("invalid_result", "CMUX artifact identities are invalid") + for item in artifacts: + entry = _exact_keys( + item, + {"class", "path_class", "sha256", "bytes"}, + "CMUX artifact identity", + ) + if ( + not isinstance(entry.get("class"), str) + or not entry["class"] + or entry.get("path_class") != "repository_output" + or not isinstance(entry.get("sha256"), str) + or DIGEST_RE.fullmatch(entry["sha256"]) is None + or type(entry.get("bytes")) is not int + or entry["bytes"] < 0 + ): + raise ContractRefusal("invalid_result", "CMUX artifact identity is invalid") + + missing = validation["missing_required_artifact_classes"] + if ( + not isinstance(missing, list) + or any(not isinstance(item, str) or not item for item in missing) + ): + raise ContractRefusal("invalid_result", "CMUX artifact validation is invalid") + observations = toolchain["observations"] + if ( + not isinstance(toolchain.get("identity"), str) + or DIGEST_RE.fullmatch(toolchain["identity"]) is None + or not isinstance(observations, dict) + or any( + not isinstance(name, str) + or not name + or not isinstance(observation, str) + for name, observation in observations.items() + ) + or toolchain["identity"] != sha256(canonical_bytes(observations)) + ): + raise ContractRefusal("invalid_result", "CMUX toolchain identity is inconsistent") + if ( + benchmark.get("state_class") not in STATE_CLASSES + or not isinstance(benchmark.get("semantic_comparison_key"), str) + or DIGEST_RE.fullmatch(benchmark["semantic_comparison_key"]) is None + or not isinstance(benchmark.get("comparison_context_key"), str) + or DIGEST_RE.fullmatch(benchmark["comparison_context_key"]) is None + ): + raise ContractRefusal("invalid_result", "CMUX benchmark identity is invalid") + semantic_key = _cmux_semantic_key(value) + if benchmark["semantic_comparison_key"] != semantic_key: + raise ContractRefusal("invalid_result", "CMUX semantic comparison identity is inconsistent") + if benchmark["comparison_context_key"] != _cmux_context_key( + semantic_key, + benchmark["state_class"], + toolchain["identity"], + ): + raise ContractRefusal("invalid_result", "CMUX comparison context is inconsistent") + if ( + cleanup.get("state") not in {"complete", "forced", "incomplete"} + or type(cleanup.get("process_group_settled")) is not bool + ): + raise ContractRefusal("invalid_result", "CMUX cleanup evidence is invalid") + if value["result"] == "ambiguous": + if cleanup != {"state": "forced", "process_group_settled": False}: + raise ContractRefusal("invalid_result", "ambiguous CMUX result lacks forced cleanup evidence") + elif cleanup != {"state": "complete", "process_group_settled": True}: + raise ContractRefusal("invalid_result", "terminal CMUX result lacks complete cleanup") + if value["result"] == "passed" and ( + value["exit_code"] != 0 or missing + ): + raise ContractRefusal("invalid_result", "passed CMUX result is inconsistent") + if value["result"] == "failed" and value["exit_code"] == 0 and not missing: + raise ContractRefusal("invalid_result", "failed CMUX result is inconsistent") + timings = value["stage_timings"] + if ( + not isinstance(timings, list) + or not timings + or any( + not isinstance(item, dict) + or set(item) != {"stage", "seconds"} + or not isinstance(item["stage"], str) + or not item["stage"] + or isinstance(item["seconds"], bool) + or not isinstance(item["seconds"], (int, float)) + or item["seconds"] < 0 + for item in timings + ) + ): + raise ContractRefusal("invalid_result", "CMUX stage timings are invalid") + if ( + not isinstance(resource.get("resource_class"), str) + or not resource["resource_class"] + or type(resource.get("cpu_count")) is not int + or resource["cpu_count"] <= 0 + or (resource.get("memory_bytes") is not None and ( + type(resource["memory_bytes"]) is not int or resource["memory_bytes"] <= 0 + )) + or not isinstance(resource.get("architecture"), str) + or not resource["architecture"] + ): + raise ContractRefusal("invalid_result", "CMUX resource summary is invalid") + return value + + +def decode_cmux_result(raw: bytes) -> dict[str, object]: + value = decode_json(raw, MAX_DOCUMENT_BYTES * 4) + if not isinstance(value, dict): + raise ContractRefusal("invalid_result", "CMUX semantic result must be an object") + canonical = canonical_bytes(value) + b"\n" + if raw != canonical: + raise ContractRefusal("invalid_result", "CMUX semantic result is not canonical") + if ( + value.get("document_type") != "cmux-workload-result" + or type(value.get("schema_version")) is not int + or value.get("schema_version") != 1 + ): + raise ContractRefusal( + "invalid_result", "CMUX semantic result contract is unsupported" + ) + return validate_cmux_result(value) + + +def observe(request: Request, raw_result: bytes) -> dict[str, object]: + result = decode_cmux_result(raw_result) + expected_source = { + "repository": request.repository, + "commit": request.commit, + "tree": request.tree, + } + expected_profile = { + "id": request.profile_id, + "generation": request.profile_generation, + } + if result.get("source") != expected_source or result.get("profile") != expected_profile: + raise ContractRefusal( + "result_mismatch", "CMUX result source/profile differs from request" + ) + if result.get("parameters") != request.parameters: + raise ContractRefusal( + "result_mismatch", "CMUX result parameters differ from request" + ) + benchmark = result.get("benchmark") + if ( + not isinstance(benchmark, dict) + or benchmark.get("state_class") != request.benchmark_state_class + ): + raise ContractRefusal( + "result_mismatch", "CMUX result benchmark state differs from request" + ) + validator = result.get("semantic_validator") + artifacts = result.get("artifact_identities") + cleanup = result.get("cleanup") + if ( + not isinstance(validator, str) + or not validator + or not isinstance(artifacts, list) + or not isinstance(cleanup, dict) + ): + raise ContractRefusal("invalid_result", "CMUX semantic evidence is incomplete") + outer: dict[str, object] = { + "document_type": OBSERVATION_TYPE, + "schema_version": SCHEMA_VERSION, + "external_request_ref": request.external_request_ref, + "request_sha256": sha256(canonical_bytes(request_document(request))), + "execution_binding_sha256": execution_binding(request), + "source": expected_source, + "profile": expected_profile, + "state": OUTER_STATES[result["result"]], + "cmux_semantic_result_sha256": sha256(raw_result), + "cmux_semantic_validator": validator, + "authority": AUTHORITY, + } + if request.work_ref is not None: + outer["correlation"] = {"work_ref": request.work_ref} + return outer + + +def write_document(value: object) -> None: + raw = canonical_bytes(value) + b"\n" + if len(raw) > MAX_DOCUMENT_BYTES: + raise ContractRefusal("invalid_document", "output exceeds its fixed ceiling") + sys.stdout.buffer.write(raw) + + +def parser() -> argparse.ArgumentParser: + root = argparse.ArgumentParser(description=__doc__) + commands = root.add_subparsers(dest="command", required=True) + commands.add_parser("plan") + observe_parser = commands.add_parser("observe") + observe_parser.add_argument("--request", required=True) + observe_parser.add_argument("--result", required=True) + return root + + +def main() -> int: + args = parser().parse_args() + try: + if args.command == "plan": + request = decode_request(sys.stdin.buffer.read(MAX_DOCUMENT_BYTES + 1)) + write_document(plan(request)) + return 0 + if args.command == "observe": + request = decode_request(Path(args.request).read_bytes()) + result = Path(args.result).read_bytes() + write_document(observe(request, result)) + return 0 + raise AssertionError(args.command) + except (ContractRefusal, OSError) as error: + code = error.code if isinstance(error, ContractRefusal) else "io_error" + print(f"cmux workload request refused [{code}]: {error}", file=sys.stderr) + return 64 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/scripts/test-cmux-workload-request.py b/scripts/test-cmux-workload-request.py new file mode 100644 index 000000000..80c6fba01 --- /dev/null +++ b/scripts/test-cmux-workload-request.py @@ -0,0 +1,295 @@ +#!/usr/bin/env python3 +from __future__ import annotations + +import copy +import importlib.util +import json +from pathlib import Path +import subprocess +import sys +import unittest + +ROOT = Path(__file__).resolve().parents[1] +MODULE_PATH = ROOT / "scripts/cmux_workload_request.py" +SPEC = importlib.util.spec_from_file_location("cmux_workload_request", MODULE_PATH) +assert SPEC and SPEC.loader +module = importlib.util.module_from_spec(SPEC) +sys.modules["cmux_workload_request"] = module +SPEC.loader.exec_module(module) + +BASE = { + "document_type": "glaeda-cmux-workload-request", + "schema_version": 1, + "external_request_ref": "cmux:workload:fixture-1", + "source": { + "repository": "manaflow-ai/cmux", + "commit": "1" * 40, + "tree": "2" * 40, + }, + "profile": {"id": "cmux.ci.guard", "generation": 1}, + "benchmark_state_class": "cold", + "reuse_hint": "no_preference", + "correlation": {"work_ref": "cmux:work:fixture-1"}, +} + + +def raw(value: object) -> bytes: + return module.canonical_bytes(value) + b"\n" + + +def cmux_result(request: module.Request) -> dict[str, object]: + observations = {"python": "3.12.0"} + result: dict[str, object] = { + "document_type": "cmux-workload-result", + "schema_version": 1, + "source": { + "repository": request.repository, + "commit": request.commit, + "tree": request.tree, + }, + "profile": { + "id": request.profile_id, + "generation": request.profile_generation, + }, + "semantic_validator": "cmux.ci-guard/v1", + "expected_result_class": "cmux.ci-guard-result/v1", + "result": "passed", + "parameters": request.parameters, + "runtime_input_identities": [], + "artifact_identities": [], + "validation": {"missing_required_artifact_classes": []}, + "stage_timings": [{"stage": "test", "seconds": 1.25}], + "resource_summary": { + "resource_class": "cmux-linux-ci-small", + "cpu_count": 4, + "memory_bytes": 17179869184, + "architecture": "x86_64", + }, + "toolchain": { + "identity": module.sha256(module.canonical_bytes(observations)), + "observations": observations, + }, + "benchmark": { + "state_class": request.benchmark_state_class, + "semantic_comparison_key": "sha256:" + "0" * 64, + "comparison_context_key": "sha256:" + "0" * 64, + }, + "network_class": "none", + "timeout_class": "portable-short", + "cleanup": {"state": "complete", "process_group_settled": True}, + "exit_code": 0, + "started_at_unix_millis": 1000, + "ended_at_unix_millis": 2250, + } + semantic = module._cmux_semantic_key(result) + result["benchmark"]["semantic_comparison_key"] = semantic + result["benchmark"]["comparison_context_key"] = module._cmux_context_key( + semantic, + request.benchmark_state_class, + result["toolchain"]["identity"], + ) + return result + + +class CmuxWorkloadRequestTests(unittest.TestCase): + def test_plan_resolves_only_fixed_repository_adapter(self) -> None: + request = module.decode_request(raw(BASE)) + plan = module.plan(request) + self.assertEqual(plan["state"], "planned") + self.assertEqual( + plan["resolved_adapter"], + { + "kind": "cmux-repository-profile/v1", + "repository_runner": "scripts/ci/cmux_workload_profile.py", + "semantic_result_contract": "cmux-workload-result/v1", + }, + ) + self.assertEqual(plan["authority"], module.AUTHORITY) + encoded = module.canonical_bytes(plan) + for forbidden in (b"argv", b"cwd", b"machine", b"resource_class", b"state_root"): + self.assertNotIn(forbidden, encoded) + + def test_execution_binding_ignores_caller_correlation_and_reuse_hint(self) -> None: + first = module.decode_request(raw(BASE)) + changed = copy.deepcopy(BASE) + changed["external_request_ref"] = "cmux:workload:other" + changed["correlation"]["work_ref"] = "cmux:work:other" + changed["reuse_hint"] = "prefer_valid_reuse" + second = module.decode_request(raw(changed)) + self.assertEqual(module.execution_binding(first), module.execution_binding(second)) + self.assertNotEqual( + module.sha256(module.canonical_bytes(module.request_document(first))), + module.sha256(module.canonical_bytes(module.request_document(second))), + ) + + def test_profile_generation_changes_execution_binding(self) -> None: + first = module.decode_request(raw(BASE)) + changed = copy.deepcopy(BASE) + changed["profile"]["generation"] = 2 + second = module.decode_request(raw(changed)) + self.assertNotEqual(module.execution_binding(first), module.execution_binding(second)) + + def test_rejects_execution_controls(self) -> None: + for field, value in { + "argv": ["caller-selected-program"], + "cwd": "/caller/path", + "machine": "caller-machine", + "backend": "caller-backend", + "resource_class": "caller-resource", + "state_root": "/caller/cache", + }.items(): + with self.subTest(field=field): + changed = copy.deepcopy(BASE) + changed[field] = value + with self.assertRaisesRegex(module.ContractRefusal, "request fields"): + module.decode_request(raw(changed)) + + def test_observation_preserves_cmux_semantic_result(self) -> None: + request = module.decode_request(raw(BASE)) + result_raw = raw(cmux_result(request)) + observation = module.observe(request, result_raw) + self.assertEqual(observation["state"], "succeeded") + self.assertEqual(observation["cmux_semantic_validator"], "cmux.ci-guard/v1") + self.assertEqual( + observation["cmux_semantic_result_sha256"], module.sha256(result_raw) + ) + + def test_observation_rejects_source_or_profile_drift(self) -> None: + request = module.decode_request(raw(BASE)) + result = cmux_result(request) + result["profile"]["generation"] = 2 + semantic = module._cmux_semantic_key(result) + result["benchmark"]["semantic_comparison_key"] = semantic + result["benchmark"]["comparison_context_key"] = module._cmux_context_key( + semantic, + result["benchmark"]["state_class"], + result["toolchain"]["identity"], + ) + with self.assertRaisesRegex(module.ContractRefusal, "differs from request"): + module.observe(request, raw(result)) + + def test_observation_rejects_incomplete_passed_result(self) -> None: + request = module.decode_request(raw(BASE)) + + missing_artifact = cmux_result(request) + missing_artifact["validation"]["missing_required_artifact_classes"] = [ + "cmux-required-product" + ] + with self.assertRaisesRegex(module.ContractRefusal, "passed CMUX result"): + module.observe(request, raw(missing_artifact)) + + unsettled = cmux_result(request) + unsettled["cleanup"] = { + "state": "forced", + "process_group_settled": False, + } + with self.assertRaisesRegex(module.ContractRefusal, "complete cleanup"): + module.observe(request, raw(unsettled)) + + bad_exit = cmux_result(request) + bad_exit["exit_code"] = 1 + with self.assertRaisesRegex(module.ContractRefusal, "passed CMUX result"): + module.observe(request, raw(bad_exit)) + + def test_observation_rejects_result_structure_and_toolchain_tampering(self) -> None: + request = module.decode_request(raw(BASE)) + + extra = cmux_result(request) + extra["unexpected_private_field"] = "must-not-pass" + with self.assertRaisesRegex(module.ContractRefusal, "result fields"): + module.observe(request, raw(extra)) + + tampered_toolchain = cmux_result(request) + tampered_toolchain["toolchain"]["observations"]["python"] = "different" + with self.assertRaisesRegex(module.ContractRefusal, "toolchain identity"): + module.observe(request, raw(tampered_toolchain)) + + missing_cpu = cmux_result(request) + missing_cpu["resource_summary"]["cpu_count"] = None + with self.assertRaisesRegex(module.ContractRefusal, "resource summary"): + module.observe(request, raw(missing_cpu)) + + def test_negative_result_requires_settled_cleanup_unless_ambiguous(self) -> None: + request = module.decode_request(raw(BASE)) + + failed = cmux_result(request) + failed["result"] = "failed" + failed["exit_code"] = 1 + observation = module.observe(request, raw(failed)) + self.assertEqual(observation["state"], "failed") + + ambiguous = cmux_result(request) + ambiguous["result"] = "ambiguous" + ambiguous["exit_code"] = -15 + ambiguous["cleanup"] = { + "state": "forced", + "process_group_settled": False, + } + observation = module.observe(request, raw(ambiguous)) + self.assertEqual(observation["state"], "ambiguous") + + contradictory = cmux_result(request) + contradictory["result"] = "ambiguous" + with self.assertRaisesRegex(module.ContractRefusal, "ambiguous CMUX result"): + module.observe(request, raw(contradictory)) + + def test_observation_rejects_self_inconsistent_cmux_result(self) -> None: + request = module.decode_request(raw(BASE)) + cases = [] + + extra = cmux_result(request) + extra["unexpected"] = True + cases.append(("fields", extra)) + + toolchain = cmux_result(request) + toolchain["toolchain"]["observations"]["python"] = "changed" + cases.append(("toolchain", toolchain)) + + semantic = cmux_result(request) + semantic["benchmark"]["semantic_comparison_key"] = "sha256:" + "f" * 64 + cases.append(("semantic comparison", semantic)) + + cleanup = cmux_result(request) + cleanup["cleanup"] = {"state": "forced", "process_group_settled": False} + cases.append(("cleanup", cleanup)) + + artifacts = cmux_result(request) + artifacts["validation"]["missing_required_artifact_classes"] = [ + "cmux.required-artifact/v1" + ] + cases.append(("passed CMUX result", artifacts)) + + for label, result in cases: + with self.subTest(label=label): + with self.assertRaises(module.ContractRefusal): + module.observe(request, raw(result)) + + def test_ambiguous_result_requires_forced_unsettled_cleanup(self) -> None: + request = module.decode_request(raw(BASE)) + result = cmux_result(request) + result["result"] = "ambiguous" + result["exit_code"] = 0 + result["cleanup"] = {"state": "forced", "process_group_settled": False} + observation = module.observe(request, raw(result)) + self.assertEqual(observation["state"], "ambiguous") + + result["cleanup"] = {"state": "complete", "process_group_settled": True} + with self.assertRaisesRegex(module.ContractRefusal, "ambiguous"): + module.observe(request, raw(result)) + + def test_repository_plan_fixture_round_trip(self) -> None: + request_path = ROOT / "docs/experiments/cmux-workload-profile/request.json" + plan_path = ROOT / "docs/experiments/cmux-workload-profile/plan.json" + planned = subprocess.run( + [sys.executable, str(MODULE_PATH), "plan"], + input=request_path.read_bytes(), + stdout=subprocess.PIPE, + stderr=subprocess.PIPE, + check=False, + ) + self.assertEqual(planned.returncode, 0, planned.stderr.decode()) + self.assertEqual(planned.stdout, plan_path.read_bytes()) + + +if __name__ == "__main__": + unittest.main()