#!/usr/bin/env python3 """Serve trusted Drone configurations and execute only their authorized job phases. /config authenticates Drone itself. /jobs// authenticates a short-lived job credential and enforces the event policy stored when that job was created. """ import hashlib import hmac import json import os import re import secrets import threading import time from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer from pathlib import Path from policy import drone_config, registration, verify_request from workflow import Workflow BUNDLE_ROOT = Path(__file__).resolve().parents[1] JOB_STORE = BUNDLE_ROOT / "config/private/ci-jobs" JOB_STORE.mkdir(mode=0o700, parents=True, exist_ok=True) JOB_LOCKS = {} LOCK_REGISTRY_MUTEX = threading.Lock() MAX_REQUEST_BYTES = 2 * 1024 * 1024 JOB_LIFETIME_SECONDS = 86400 JOB_ROUTE = r"/jobs/([0-9]{14}-[a-f0-9]{20})/([a-z-]+)" def job_lock(job_id): """Serialize phase execution and state writes for a single job.""" with LOCK_REGISTRY_MUTEX: return JOB_LOCKS.setdefault(job_id, threading.Lock()) class Handler(BaseHTTPRequestHandler): """Expose health, configuration, and phase endpoints on loopback only.""" def log_message(self, *args): # The default HTTP logger could expose job URLs or credentials. pass def answer(self, status, data): """Send one JSON response with an explicit content length.""" body = json.dumps(data).encode() self.send_response(status) self.send_header("Content-Type", "application/json") self.send_header("Content-Length", str(len(body))) self.end_headers() self.wfile.write(body) def do_GET(self): if self.path == "/healthz": self.answer(200, {"status": "ok"}) else: self.answer(404, {"error": "Not found"}) def do_POST(self): """Validate request size, dispatch, and translate errors into JSON.""" try: length = int(self.headers.get("Content-Length", "0")) if not 0 < length < MAX_REQUEST_BYTES: raise ValueError("Invalid request size") body = self.rfile.read(length) if self.path == "/config": self.create_job_configuration(body) return match = re.fullmatch(JOB_ROUTE, self.path) if not match: raise ValueError("Unknown endpoint") job_id, phase = match.groups() self.execute_job_phase(job_id, phase) except (ValueError, KeyError, FileNotFoundError) as error: self.answer(403, {"error": str(error)}) except Exception as error: # Keep unexpected error details and credentials out of client responses. print("CI service error:", type(error).__name__, flush=True) self.answer(500, {"error": "CI service error; inspect service logs"}) def create_job_configuration(self, body): """Authenticate Drone and bind a new credential to the permitted phases.""" verify_request(self.headers, body, os.environ["DRONE_YAML_SECRET"]) request = json.loads(body) allowed_repositories = os.environ["CI_REPOSITORIES"].split(",") job = registration(request, allowed_repositories) job["id"] = time.strftime("%Y%m%d%H%M%S") + "-" + secrets.token_hex(10) token = secrets.token_urlsafe(32) job["token_hash"] = hashlib.sha256(token.encode()).hexdigest() job["created"] = time.time() job["results"] = {} (JOB_STORE / (job["id"] + ".json")).write_text(json.dumps(job)) configuration = drone_config(job, os.environ["CI_ENDPOINT"], token) self.answer(200, {"data": configuration}) def execute_job_phase(self, job_id, phase): """Check stored authorization before executing or returning a phase result.""" with job_lock(job_id): state_path = JOB_STORE / (job_id + ".json") job = json.loads(state_path.read_text()) supplied_token = self.headers.get("Authorization", "").removeprefix( "Bearer " ) supplied_hash = hashlib.sha256(supplied_token.encode()).hexdigest() if not hmac.compare_digest(supplied_hash, job["token_hash"]): raise ValueError("Invalid job token") if time.time() - job["created"] > JOB_LIFETIME_SECONDS: raise ValueError("Job expired") if phase not in job["phases"]: raise ValueError("This phase is forbidden for this event") workflow = Workflow(job) # Completed actions are idempotent: retries return the recorded result. if phase not in job["results"]: job["results"][phase] = workflow.execute(phase) temporary_path = state_path.with_suffix(".tmp") temporary_path.write_text(json.dumps(job)) temporary_path.replace(state_path) self.answer( 200, { **job["results"][phase], "log": workflow.log(phase), "reports": str(workflow.run_dir / "reports"), }, ) if __name__ == "__main__": ThreadingHTTPServer( ("127.0.0.1", int(os.environ.get("CI_PORT", "19095"))), Handler ).serve_forever()