"""Host-controlled phase implementations. Repository files are inputs, never CI commands.""" import base64 import hashlib import http.client import json import os import re import shutil import subprocess import time import urllib.parse import urllib.request import urllib.error from pathlib import Path from policy import SCANS, RELEASE import artifacts BUNDLE_ROOT = Path(__file__).resolve().parents[1] HOST_BUNDLE_ROOT = Path(os.environ.get("HOST_BUNDLE_ROOT", BUNDLE_ROOT)) ENGINE = os.environ.get("CONTAINER_ENGINE", "docker") IMAGES = { **artifacts.IMAGES, "sonar": "docker.io/sonarsource/sonar-scanner-cli:12.1.0.3233_8.0.1", } def host(path): """Translate a bundle path to the host path needed for container bind mounts.""" return str(HOST_BUNDLE_ROOT / path.relative_to(BUNDLE_ROOT)) def call( url, method="GET", values=None, token=None, basic=None, file=None, dest=None, openqa=False, ): """Call a configured loopback service, optionally streaming an artifact file.""" parsed_url = urllib.parse.urlsplit(url) if parsed_url.scheme != "http" or parsed_url.hostname not in ( "localhost", "127.0.0.1", ): raise RuntimeError("Only configured localhost service endpoints are supported") headers = {} data = None if token: headers["Authorization"] = "Bearer " + token if basic: headers["Authorization"] = "Basic " + base64.b64encode(basic.encode()).decode() if openqa: headers["X-API-Microtime"] = str(time.time()) if values is not None: data = urllib.parse.urlencode(values).encode() headers["Content-Type"] = "application/x-www-form-urlencoded" if file: data = file.open("rb") headers["Content-Length"] = str(file.stat().st_size) headers["Content-Type"] = "application/octet-stream" connection = http.client.HTTPConnection( parsed_url.hostname, parsed_url.port, timeout=600 ) try: connection.request( method, parsed_url.path + ("?" + parsed_url.query if parsed_url.query else ""), data, headers, ) response = connection.getresponse() if response.status >= 300: raise RuntimeError( "Service request failed: HTTP " + str(response.status) + " " + parsed_url.path ) if dest: with dest.open("wb") as file_handle: shutil.copyfileobj(response, file_handle) return None raw = response.read() try: return json.loads(raw) if raw else {} except ValueError: return {"text": raw.decode(errors="replace")} finally: connection.close() if file: data.close() class Workflow: """Execute trusted phases and retain per-job artifact directories and reports.""" def __init__(self, job): self.job = job self.run_dir = BUNDLE_ROOT / "runs" / job["id"] self.source_dir = self.run_dir / "source" self.candidate_dir = self.run_dir / "candidate" self.tested_dir = self.run_dir / "tested" self.reports_dir = self.run_dir / "reports" for d in [self.run_dir, self.candidate_dir, self.tested_dir, self.reports_dir]: d.mkdir(parents=True, exist_ok=True) self.container = "nd-candidate-" + job["id"] self.stream = None def log(self, phase): path = self.reports_dir / (phase + ".log") return path.read_text(errors="replace")[-60000:] if path.exists() else "" def command(self, args, **kwargs): result = subprocess.run( [str(x) for x in args], stdout=self.stream, stderr=subprocess.STDOUT, **kwargs, ) if result.returncode: raise RuntimeError( "Command failed with exit code " + str(result.returncode) + "; see this phase log" ) def tool(self, name, args, network="none", extra=()): self.command( [ ENGINE, "run", "--rm", "--pull=never", "--network", network, "--user", "0", "-v", host(self.source_dir) + ":/src:ro,z", "-v", host(self.candidate_dir) + ":/out:z", "-v", host(self.reports_dir) + ":/reports:z", *extra, IMAGES[name], *args, ] ) def execute(self, phase): """Enforce prerequisites, run one phase, and record a passed or failed result.""" with (self.reports_dir / (phase + ".log")).open("w") as self.stream: try: results = self.job["results"] if ( phase in SCANS and results.get("checkout", {}).get("status") != "passed" ): raise RuntimeError("Checkout did not pass") if phase in RELEASE: if not self.job["release_allowed"]: raise RuntimeError( "Release phases are forbidden for this event" ) previous = ["security-gate", *RELEASE[: RELEASE.index(phase)]] if any( results.get(path, {}).get("status") != "passed" for path in previous ): raise RuntimeError("An earlier required phase did not pass") getattr(self, phase.replace("-", "_"))() result = {"status": "passed"} except Exception as error: # Do not echo subprocess arguments or service credentials. self.stream.write("FAILED: " + str(error) + "\n") result = {"status": "failed"} self.stream.flush() if phase in SCANS or phase == "security-gate": self.write_summary({**self.job["results"], phase: result}) return result def write_summary(self, results): (self.reports_dir / "gates.json").write_text( json.dumps( { "commit": self.job["commit"], "event": self.job["event"], "ref": self.job["ref"], "release_allowed": self.job["release_allowed"], "gates": results, }, indent=2, ) + "\n" ) def checkout(self): """Fetch full history and check out the exact commit authorized for this job.""" if self.source_dir.exists(): raise RuntimeError("Checkout directory already exists; start a new build") askpass_path = self.run_dir / "askpass.sh" askpass_path.write_text( '#!/bin/sh\ncase "$1" in *Username*) printf "%s\\n" "$GIT_USERNAME";; *) printf "%s\\n" "$GIT_PASSWORD";; esac\n' ) askpass_path.chmod(0o700) environment = { **os.environ, "GIT_ASKPASS": str(askpass_path), "GIT_TERMINAL_PROMPT": "0", } try: self.command( [ "git", "clone", "--no-checkout", "http://localhost:" + os.environ["GITEA_PORT"] + "/" + self.job["repo"] + ".git", self.source_dir, ], env=environment, ) self.command( ["git", "-C", self.source_dir, "fetch", "origin", self.job["ref"]], env=environment, ) self.command( [ "git", "-C", self.source_dir, "checkout", "--detach", self.job["commit"], ], env=environment, ) finally: askpass_path.unlink(missing_ok=True) actual = subprocess.check_output( ["git", "-C", str(self.source_dir), "rev-parse", "HEAD"], text=True ).strip() if actual != self.job["commit"]: raise RuntimeError("Checkout commit mismatch") for path in self.source_dir.rglob("*"): if path.is_symlink(): raise RuntimeError("Symlinks are not supported in this demo checkout") self.stream.write("Checked out " + actual + "\n") def gitleaks(self): """Scan files and history with the host-managed secret detection policy.""" self.tool( "gitleaks", ["bash", "/managed-gitleaks.sh"], extra=[ "-v", host(BUNDLE_ROOT / "ci/gitleaks.sh") + ":/managed-gitleaks.sh:ro,z", ], ) def semgrep(self): """Apply the bundled rules while ignoring repository suppression settings.""" self.tool( "semgrep", [ "semgrep", "scan", "--config", "/rules.yml", "--error", "--strict", "--metrics=off", "--disable-version-check", "--disable-nosem", "--no-git-ignore", "--x-ignore-semgrepignore-files", "--exclude", ".git", "--json", "--output", "/reports/semgrep.json", "/src", ], extra=["-v", host(BUNDLE_ROOT / "config/semgrep.yml") + ":/rules.yml:ro,z"], ) def sonarqube(self): """Analyze an isolated project, check its gate, and revoke its scoped token.""" base = os.environ["SONAR_URL"] admin = "admin:" + os.environ["SONAR_PASSWORD"] key = "nd-" + self.job["id"] token = None call( base + "/api/projects/create", "POST", {"project": key, "name": self.job["repo"] + " " + self.job["commit"][:12]}, basic=admin, ) call( base + "/api/qualitygates/select", "POST", {"gateName": "Offline DevSecOps Source v1", "projectKey": key}, basic=admin, ) token_name = "ci-" + self.job["id"] token = call( base + "/api/user_tokens/generate", "POST", {"name": token_name, "type": "PROJECT_ANALYSIS_TOKEN", "projectKey": key}, basic=admin, )["token"] directory = self.run_dir / "sonar-work" directory.mkdir() directory.chmod(0o777) try: self.tool( "sonar", [ "sonar-scanner", "-Dproject.settings=/managed-sonar.properties", "-Dsonar.projectKey=" + key, "-Dsonar.projectBaseDir=/src", "-Dsonar.sources=.", "-Dsonar.working.directory=/sonar-work", "-Dsonar.scm.revision=" + self.job["commit"], "-Dsonar.scanner.skipJreProvisioning=true", ], network="host", extra=[ "-v", host(BUNDLE_ROOT / "config/sonar-project.properties") + ":/managed-sonar.properties:ro,z", "-e", "SONAR_HOST_URL=" + base, "-e", "SONAR_TOKEN=" + token, "-v", host(directory) + ":/sonar-work:z", ], ) task_properties = dict( line.split("=", 1) for line in (directory / "report-task.txt").read_text().splitlines() if "=" in line ) task_id = task_properties["ceTaskId"] if not re.fullmatch(r"[A-Za-z0-9_-]+", task_id): raise RuntimeError("Invalid Sonar task ID") for _ in range(180): task = call(base + "/api/ce/task?id=" + task_id, basic=admin)["task"] if task["status"] in ("SUCCESS", "FAILED", "CANCELED"): break time.sleep(2) else: raise RuntimeError("Sonar analysis timed out") if task["status"] != "SUCCESS": raise RuntimeError("Sonar compute task did not succeed") gate = call( base + "/api/qualitygates/project_status?analysisId=" + urllib.parse.quote(task["analysisId"]), basic=admin, ) (self.reports_dir / "sonarqube.json").write_text( json.dumps( {"project": key, "task": task, "quality_gate": gate}, indent=2 ) ) if gate["projectStatus"]["status"] != "OK": raise RuntimeError("SonarQube quality gate failed") self.stream.write( "SonarQube quality gate OK for this analysis. Project: " + key + "\n" ) finally: call( base + "/api/user_tokens/revoke", "POST", {"name": token_name}, basic=admin, ) def security_gate(self): """Require every source scanner to pass before any artifact phase can run.""" failed = [ path for path in SCANS if self.job["results"].get(path, {}).get("status") != "passed" ] if failed: raise RuntimeError("Source security gate blocked: " + ", ".join(failed)) self.stream.write("All three source checks passed.\n") def os_build(self): """Build once without network access and export the resulting container image.""" tag = "localhost/newdevsecops/app:" + self.job["id"] self.command( [ ENGINE, "build", "--network=none", "--pull=" + ("never" if ENGINE == "podman" else "false"), "-t", tag, self.source_dir, ], env={**os.environ, "DOCKER_BUILDKIT": "0"}, ) self.command( [ ENGINE, "save", *(["--format", "docker-archive"] if ENGINE == "podman" else []), "-o", self.candidate_dir / "app.tar", tag, ] ) info = json.loads(subprocess.check_output([ENGINE, "image", "inspect", tag]))[0] (self.candidate_dir / "build.json").write_text( json.dumps( { "build_id": self.job["id"], "commit": self.job["commit"], "image_id": info["Id"], "image_tag": tag, "artifact_kind": "container-image", }, indent=2, ) ) def syft_sbom(self): """Generate both SBOM formats from the exported image archive.""" self.tool("syft", ["tool-security", "image", "/out/app.tar", "/out"]) def vulnerability_gate(self): """Run both vulnerability scanners and fail if either rejects the artifact.""" failures = [] for tool, args in [ ( "grype", ["tool-security", "sbom", "/out/sbom.syft.json", "/reports/grype"], ), ("trivy", ["tool-security", "image", "/out/app.tar", "/reports/trivy"]), ]: try: self.tool(tool, args) except Exception: failures.append(tool) if failures: raise RuntimeError("Vulnerability checks failed: " + ", ".join(failures)) def cosign_sign(self): """Sign a manifest binding artifact bytes and SBOMs to this commit.""" manifest = { "build_id": self.job["id"], "commit": self.job["commit"], "files": { n: artifacts.sha(self.candidate_dir / n) for n in artifacts.FILES }, } (self.candidate_dir / "manifest.json").write_text( json.dumps(manifest, sort_keys=True, indent=2) ) self.command( [ BUNDLE_ROOT / "bin/cosign", "sign-blob", "--use-signing-config=false", "--tlog-upload=false", "--key", BUNDLE_ROOT / "config/private/cosign.key", "--bundle", self.candidate_dir / "manifest.sigstore.json", self.candidate_dir / "manifest.json", ] ) artifacts.verify_candidate( self.candidate_dir, BUNDLE_ROOT / "config/cosign.pub" ) def jfrog_candidate(self): """Upload the signed candidate, download it again, and verify its provenance.""" token = os.environ.get("JFROG_TOKEN", "") if not token: raise RuntimeError( "Configure the Artifactory license and JFROG_TOKEN; publication is blocked" ) base = ( os.environ["JFROG_URL"] + "/" + os.environ["JFROG_CANDIDATE_REPO"] + "/my-app/" + self.job["id"] ) for name in [*artifacts.FILES, "manifest.json", "manifest.sigstore.json"]: call(base + "/" + name, "PUT", token=token, file=self.candidate_dir / name) for name in [*artifacts.FILES, "manifest.json", "manifest.sigstore.json"]: call(base + "/" + name, token=token, dest=self.tested_dir / name) manifest = artifacts.verify_candidate( self.tested_dir, BUNDLE_ROOT / "config/cosign.pub" ) if ( manifest["commit"] != self.job["commit"] or manifest["build_id"] != self.job["id"] ): raise RuntimeError("Downloaded artifact provenance mismatch") def openqa(self): """Run regression tests against the downloaded image and verify job identity.""" self.command([ENGINE, "load", "-i", self.tested_dir / "app.tar"]) identity = json.loads((self.tested_dir / "build.json").read_text())["image_id"] self.command( [ ENGINE, "run", "-d", "--pull=never", "--name", self.container, "--network", "newdevsecops", "--read-only", "--cap-drop=ALL", "--security-opt=no-new-privileges", identity, ] ) base = os.environ["OPENQA_URL"] + "/api/v1/" digest = artifacts.sha(self.tested_dir / "app.tar") created = call( base + "jobs", "POST", { "TEST": "http-regression", "DISTRI": "devsecops", "VERSION": "1", "ARCH": "x86_64", "MACHINE": "null", "BACKEND": "null", "BUILD": digest, "CASEDIR": "/opt/devsecops-tests", "NEEDLES_DIR": "/opt/devsecops-tests/needles", "DEVSECOPS_COMMIT": self.job["commit"], "DEVSECOPS_TARGET": "http://" + self.container + ":8080", "WORKER_CLASS": "qemu_x86_64", "WORKER_MAX_JOB_TIME": "900", }, openqa=True, ) (self.reports_dir / "openqa-job.json").write_text(json.dumps(created)) for _ in range(300): job = call(base + "jobs/" + str(created["id"]), openqa=True)["job"] if job["state"] in ("done", "cancelled"): (self.reports_dir / "openqa.json").write_text(json.dumps(job, indent=2)) if ( job["result"] != "passed" or job["settings"]["BUILD"] != digest or job["settings"]["DEVSECOPS_COMMIT"] != self.job["commit"] ): raise RuntimeError("openQA did not pass for this artifact") return time.sleep(3) raise RuntimeError("openQA job timed out") def security_regression(self): """Actively scan the isolated candidate with ZAP; nonzero results block release.""" directory = self.reports_dir / "zap" directory.mkdir() directory.chmod(0o777) self.command( [ ENGINE, "run", "--rm", "--pull=never", "--network", "newdevsecops", "-v", host(directory) + ":/zap/wrk:z", IMAGES["zap"], "zap-full-scan.py", "-t", "http://" + self.container + ":8080", "-m", "1", "-T", "5", "-r", "zap.html", "-J", "zap.json", "-z", "-config autoupdate.checkOnStart=false -config autoupdate.checkAddonUpdates=false", ] ) def release(self): """Publish verified signed evidence, uploading the release marker last.""" gates = {k: v["status"] for k, v in self.job["results"].items()} local = self.run_dir / "release-staging" artifacts.publish_release( self.tested_dir, self.reports_dir, local, {"build_id": self.job["id"], "commit": self.job["commit"], "gates": gates}, BUNDLE_ROOT / "config/private/cosign.key", BUNDLE_ROOT / "config/cosign.pub", ) base = ( os.environ["JFROG_URL"] + "/" + os.environ["JFROG_RELEASE_REPO"] + "/my-app/" + self.job["id"] ) for name in [ *artifacts.FILES, "manifest.json", "manifest.sigstore.json", "evidence.tar.gz", "RELEASE.sigstore.json", "RELEASE.json", ]: call( base + "/" + name, "PUT", token=os.environ["JFROG_TOKEN"], file=local / name, ) destination = BUNDLE_ROOT / "releases" / self.job["id"] destination.parent.mkdir(exist_ok=True) local.rename(destination) self.stream.write("Released signed artifact to " + base + "\n") def report(self): """Write the final gate summary and remove the temporary candidate container.""" self.write_summary(self.job["results"]) subprocess.run( [ENGINE, "rm", "-f", self.container], stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL, ) self.stream.write((self.reports_dir / "gates.json").read_text() + "\n")