diff --git a/bootstrap/pyproject.toml b/bootstrap/pyproject.toml index 45cb2b6a..0c3866a9 100644 --- a/bootstrap/pyproject.toml +++ b/bootstrap/pyproject.toml @@ -4,7 +4,7 @@ build-backend = "setuptools.build_meta" [project] name = "praktika-controller" -version = "0.1.7" +version = "0.1.8" description = "Thin controller for launching versioned Praktika workloads" requires-python = ">=3.10" dependencies = [ diff --git a/bootstrap/src/praktika_controller/common.py b/bootstrap/src/praktika_controller/common.py index 9aabea85..71e16b69 100644 --- a/bootstrap/src/praktika_controller/common.py +++ b/bootstrap/src/praktika_controller/common.py @@ -1,6 +1,7 @@ from __future__ import annotations from collections.abc import Callable +import contextlib import errno import importlib.util import json @@ -27,6 +28,22 @@ REPO_SUBDIR = "repo" +@contextlib.contextmanager +def profile_step(log, label): + """Time a bootstrap step and log its duration as a ``[profile]`` line. + + The controller bootstrap (clone, ephemeral merge, snapshot pack/upload, + runtime install) can take minutes on large trees (e.g. ClickHouse); wrapping + each step gives a per-step breakdown to profile against without external + tooling. A no-op if ``log`` is None.""" + start = time.time() + try: + yield + finally: + if log is not None: + log.info("[profile] %s: %.2fs", label, time.time() - start) + + class LogRateLimiter: def __init__( self, @@ -186,6 +203,59 @@ def get_github_token(region: str = "") -> str: return token +def _strip_jsonc(text: str) -> str: + """Return ``text`` with JSONC niceties removed so plain ``json.loads`` accepts + it: ``//`` line comments, ``/* */`` block comments, and trailing commas. + + Lets an operator comment out a line/field in the ci_config parameter without + silently disabling the WHOLE config (strict ``json.loads`` would reject the + comment and the caller would fall back to ``{}``). The scan is string-aware, so + a ``//`` inside a value (e.g. an ``https://`` wheel URL) or a quoted ``/*`` is + left untouched — only comments outside of strings are dropped.""" + out = [] + i, n = 0, len(text) + in_str = escaped = False + while i < n: + c = text[i] + if in_str: + out.append(c) + if escaped: + escaped = False + elif c == "\\": + escaped = True + elif c == '"': + in_str = False + i += 1 + elif c == '"': + in_str = True + out.append(c) + i += 1 + elif c == "/" and i + 1 < n and text[i + 1] == "/": + while i < n and text[i] != "\n": + i += 1 # skip to end of line + elif c == "/" and i + 1 < n and text[i + 1] == "*": + i += 2 + while i + 1 < n and not (text[i] == "*" and text[i + 1] == "/"): + i += 1 + i += 2 # skip closing */ + elif c in "}]": + # Drop a trailing comma before this closing bracket. Done here (not via + # regex on the whole text) so it stays string-aware: a "," inside a + # value like "x,]" is never touched because we only reach this branch + # outside of strings. + j = len(out) - 1 + while j >= 0 and out[j].isspace(): + j -= 1 + if j >= 0 and out[j] == ",": + del out[j] + out.append(c) + i += 1 + else: + out.append(c) + i += 1 + return "".join(out) + + def load_ci_config(region: str = "", log=None) -> dict: """Per-project, CI-wide config read from the SSM parameter ``{PRAKTIKA_PROJECT_SLUG}-ci-config`` as a JSON object. @@ -201,6 +271,10 @@ def load_ci_config(region: str = "", log=None) -> dict: ``{}`` (every feature off), so the parameter is entirely optional and absence is safe. Read fresh by the caller each run — one ``GetParameter`` is trivial next to a full CI run and keeps the config live-editable. + + JSONC is tolerated (``//`` and ``/* */`` comments, trailing commas) so an + operator can comment out a field without accidentally invalidating the whole + parameter (which would silently read as ``{}``). See ``_strip_jsonc``. """ import boto3 @@ -218,7 +292,7 @@ def load_ci_config(region: str = "", log=None) -> dict: try: ssm = boto3.client("ssm", region_name=region) value = ssm.get_parameter(Name=name)["Parameter"]["Value"] - data = json.loads(value) + data = json.loads(_strip_jsonc(value)) if not isinstance(data, dict): if log: log.warning("Controller config %s is not a JSON object; ignoring", name) @@ -759,8 +833,22 @@ def clone_repo( # reachable-SHA fetch (GitHub advertises PR commits, so a fork PR's head_sha # is fetchable from the base repo); if the sha was force-pushed away and is no # longer reachable, the fetch fails and the caller aborts — the safe outcome. - git(["fetch", "--depth=1", "origin", head_sha], cwd=clone_dir) - git(["checkout", head_sha], cwd=clone_dir) + # + # --no-tags: don't auto-follow the remote's tags (thousands on large repos + # like ClickHouse) — we only need this one sha. --no-recurse-submodules: the + # bootstrap tree never needs submodules; jobs that do fetch them themselves. + with profile_step(log, "clone: fetch head (depth 1)"): + git( + [ + "fetch", "--depth=1", "--no-tags", "--no-recurse-submodules", + "origin", head_sha, + ], + cwd=clone_dir, + ) + # Parallel checkout (checkout.workers=0 => one worker per core) speeds up + # writing a large working tree (ClickHouse is ~30k+ files) to disk. + with profile_step(log, "clone: checkout head"): + git(["-c", "checkout.workers=0", "checkout", head_sha], cwd=clone_dir) actual_sha = git(["rev-parse", "HEAD"], cwd=clone_dir).strip() if log is not None: @@ -818,27 +906,30 @@ def restore_repo_snapshot( archive_path = os.path.join(work_dir, f".repo_snapshot_{snapshot_sha[:12]}.tar.zst") try: - s3_client.download_file(bucket, key, archive_path) - - h = hashlib.sha256() - with open(archive_path, "rb") as f: - for chunk in iter(lambda: f.read(1024 * 1024), b""): - h.update(chunk) - actual_hash = h.hexdigest() - if actual_hash != expected_hash: - raise RuntimeError( - f"Repo snapshot hash mismatch for {key}: expected " - f"{expected_hash}, got {actual_hash} — refusing to run" - ) + with profile_step(log, "restore: download snapshot"): + s3_client.download_file(bucket, key, archive_path) + + with profile_step(log, "restore: verify sha256"): + h = hashlib.sha256() + with open(archive_path, "rb") as f: + for chunk in iter(lambda: f.read(1024 * 1024), b""): + h.update(chunk) + actual_hash = h.hexdigest() + if actual_hash != expected_hash: + raise RuntimeError( + f"Repo snapshot hash mismatch for {key}: expected " + f"{expected_hash}, got {actual_hash} — refusing to run" + ) # Self-contained shallow repo (.git + worktree at snapshot_sha, no history). - subprocess.run( - f"zstd -dc {shlex.quote(archive_path)} | " - f"tar -xf - -C {shlex.quote(clone_dir)}", - shell=True, - check=True, - executable="/bin/bash", - ) + with profile_step(log, "restore: unpack (zstd | tar)"): + subprocess.run( + f"zstd -dc {shlex.quote(archive_path)} | " + f"tar -xf - -C {shlex.quote(clone_dir)}", + shell=True, + check=True, + executable="/bin/bash", + ) finally: if os.path.exists(archive_path): os.remove(archive_path) diff --git a/bootstrap/src/praktika_controller/controller.py b/bootstrap/src/praktika_controller/controller.py index 284dd9f5..5867f080 100644 --- a/bootstrap/src/praktika_controller/controller.py +++ b/bootstrap/src/praktika_controller/controller.py @@ -24,6 +24,7 @@ load_ci_config, post_early_check, instance_tag, + profile_step, resolve_praktika_base_venv, restore_repo_snapshot, TaskLogCapture, @@ -36,8 +37,10 @@ prepare_repo_snapshot, read_repo_settings, ) +from praktika_controller.self_update import maybe_self_update from praktika_controller.venv_manager import ( ensure_praktika_runtime, + is_passthrough, praktika_command, venv_env, ) @@ -104,22 +107,65 @@ def _role_config(role: str) -> tuple[str, str]: raise AssertionError(f"Unhandled role: {role}") -def _resolve_runtime_source(clone_dir: str, log): - """Optional per-pool Praktika runtime source, carried on the instance's - ``praktika_runtime_source`` tag (set from the pool's ``ext['runtime_source']``). - - When set, the controller does NOT use the Praktika baked into the AMI; it - installs Praktika from this source on every task, so the pool always runs the - current checkout. The value is a filesystem path: an absolute path on the - instance, or a path relative to the cloned repo (so a tag of ``.`` installs - Praktika from the checked-out repo itself). Returns ``None`` when no runtime - source tag is set (the baked venv is used). +def _ci_config_for_message(role: str, payload, log) -> dict: + """The ci_config snapshot for this message, read **once**. + + The orchestrator (workflow role) reads it from SSM for a fresh run — the same + single read that drives both the controller self-update decision and the value + frozen into run metadata by ``handle_workflow`` (pass this snapshot to it, do + NOT re-read SSM, or a mid-message parameter change could make the orchestrator + run one controller version while freezing another into the run). A rerun reuses + the frozen copy on the event; a job runner reads the frozen copy from the task + (never SSM). See ci-config.md.""" + if not isinstance(payload, dict): + return {} + if role == ROLE_WORKFLOW and payload.get("type") != "rerun": + return load_ci_config(region=REGION, log=log) + return payload.get("ci_config") or {} + + +def _controller_version_pin(ci_config) -> str: + value = (ci_config or {}).get("praktika_controller_version", "") + return value.strip() if isinstance(value, str) else "" + + +def _resolve_runtime_source(clone_dir: str, log, ci_config=None): + """Optional Praktika runtime source override. Two sources, pin first: + + 1. ``ci_config['praktika_version']`` — a per-run version pin read once from + SSM by the orchestrator and frozen into run metadata (see ci-config.md), so + every job runner installs the SAME Praktika from run metadata, not SSM. It + supports the three pip forms: a version spec (``praktika==0.1.9``), a wheel + URL (``https://…whl``), or a repo/filesystem path (``.``). It takes + precedence over the per-pool tag below — an operator pinning a version wins + over a dev pool's checkout tag. + 2. The per-pool ``praktika_runtime_source`` instance tag (set from the pool's + ``ext['runtime_source']``): a filesystem path installed on every task so + the pool always runs the current checkout. + + When set (by either), the controller does NOT use the Praktika baked into the + AMI; it installs Praktika from the resolved source on every task. Relative + filesystem paths resolve against the cloned repo (so ``.`` installs Praktika + from the checked-out repo itself); URLs / version specs are handed to pip + verbatim. Returns ``None`` when neither is set (the baked venv is used). A failure to *read* the tag (transient IMDS/metadata error) is NOT treated as "unset": it propagates, so a pool configured to test the checkout fails the task (and the infra retry re-runs it) instead of silently passing on the baked Praktika. A genuinely absent tag returns "" from ``instance_tag`` (404) and is handled as unset below.""" + pin = (ci_config or {}).get("praktika_version", "") + pin = pin.strip() if isinstance(pin, str) else "" + if pin: + if is_passthrough(pin): + # Wheel URL or version spec: a pip target in its own right. + log.info("Praktika runtime source: ci_config pin %s", pin) + return pin + # Repo-path form: resolve like the per-pool tag (relative to the clone). + resolved = pin if os.path.isabs(pin) else os.path.join(clone_dir, pin) + log.info("Praktika runtime source: ci_config pin path %s", resolved) + return resolved + source = instance_tag("praktika_runtime_source") source = (source or "").strip() if not source: @@ -132,13 +178,13 @@ def _resolve_runtime_source(clone_dir: str, log): ) if not os.path.isabs(source): source = os.path.join(clone_dir, source) - log.info("Installing Praktika at runtime from per-pool source %s", source) + log.info("Praktika runtime source: per-pool source %s", source) return source -def _resolve_runtime(clone_dir: str, log): +def _resolve_runtime(clone_dir: str, log, ci_config=None): base_venv = resolve_praktika_base_venv(clone_dir, log) - source = _resolve_runtime_source(clone_dir, log) + source = _resolve_runtime_source(clone_dir, log, ci_config=ci_config) venv_dir = ensure_praktika_runtime( source, base_venv=base_venv, @@ -292,7 +338,9 @@ def _prepare_runner_for_task(role: str, log) -> str: return f"workdir cleanup failed: {type(e).__name__}: {e}" -def handle_workflow(event, log, queue_name: str, receive_count: int = 1): +def handle_workflow( + event, log, queue_name: str, receive_count: int = 1, ci_config=None +): wf_type = event.get("type", "unknown") log.info("Processing: %s", wf_type) @@ -345,19 +393,21 @@ def handle_workflow(event, log, queue_name: str, receive_count: int = 1): # merge) or inherited from run state on a resume. Handed to the # orchestrator so it seeds run state before dispatching any job. snapshot = None - # Out-of-repo CI config. A fresh run reads it from SSM ONCE here; a resume - # reuses the ORIGINAL run's frozen config carried on the event (from - # state.json) instead of re-reading SSM, so the resume never makes a - # different force_merge_commit / merge decision than the run it resumes - # (which would mismatch the DAG/runtime against the restored snapshot on a - # fresh-base rerun). It is reused for the force_merge_commit merge gate - # (read_repo_settings) and handed to the orchestrator (PRAKTIKA_CI_CONFIG), + # Out-of-repo CI config. The poll loop already read it ONCE for this message + # (SSM for a fresh run, the frozen copy on the event for a resume) and passed + # it in, so the self-update decision and the value frozen into run metadata + # come from the SAME snapshot — never re-read SSM here (a mid-message change + # would split them). Fall back to reading it only if a caller didn't provide + # one (e.g. a direct/test call). It gates force_merge_commit + # (read_repo_settings) and is handed to the orchestrator (PRAKTIKA_CI_CONFIG), # which freezes it into run state so every job reads the same values. See # ci-config.md. - if is_resume: - ci_config = event.get("ci_config") or {} - else: - ci_config = load_ci_config(region=REGION, log=log) + if ci_config is None: + ci_config = ( + (event.get("ci_config") or {}) + if is_resume + else load_ci_config(region=REGION, log=log) + ) resume_snapshot_key = event.get("repo_snapshot_key", "") if is_resume else "" resume_snapshot_sha = event.get("snapshot_sha", "") if is_resume else "" # Fresh-base resume (per-job "Rerun w/ fresh base" button): re-run the @@ -436,15 +486,16 @@ def _restore_original(): # both match the original merge rather than the current head. clone_dir, actual_sha, snapshot = _restore_original() else: - clone_dir, actual_sha = clone_repo( - repo, - head_sha, - pr_number, - gh_token, - work_dir=WORK_DIR, - branch=branch, - log=log, - ) + with profile_step(log, "clone: total"): + clone_dir, actual_sha = clone_repo( + repo, + head_sha, + pr_number, + gh_token, + work_dir=WORK_DIR, + branch=branch, + log=log, + ) # Stale-head guard (TOCTOU): clone_repo fetches the live # refs/pull/N/head, which can have advanced since the lambda verified @@ -510,7 +561,8 @@ def _restore_original(): "pr": pr_number, } - base_venv, venv_dir = _resolve_runtime(clone_dir, log) + with profile_step(log, "runtime: resolve venv + install praktika"): + base_venv, venv_dir = _resolve_runtime(clone_dir, log, ci_config=ci_config) event_file = os.path.join(clone_dir, "ci", "tmp", "event.json") os.makedirs(os.path.dirname(event_file), exist_ok=True) @@ -664,15 +716,16 @@ def handle_task(task, log, queue_name: str, receive_count: int = 1): # authentication so a GitHub token/login outage doesn't fail a job # that a valid snapshot could satisfy. Jobs that need `gh` authenticate # themselves via GHAuth (enable_gh_auth), independent of this. - clone_dir, actual_sha = restore_repo_snapshot( - s3, - repo_snapshot_key, - snapshot_sha, - pr_number, - work_dir=WORK_DIR, - branch=task.get("head_ref", ""), - log=log, - ) + with profile_step(log, "restore: download + verify + unpack snapshot"): + clone_dir, actual_sha = restore_repo_snapshot( + s3, + repo_snapshot_key, + snapshot_sha, + pr_number, + work_dir=WORK_DIR, + branch=task.get("head_ref", ""), + log=log, + ) else: # Cloning the head needs an authenticated remote. gh_token = get_github_token(REGION) @@ -694,7 +747,13 @@ def handle_task(task, log, queue_name: str, receive_count: int = 1): if cm_heartbeat is not None: cm_heartbeat.update(phase="resolving_runtime") - base_venv, venv_dir = _resolve_runtime(clone_dir, log) + with profile_step(log, "runtime: resolve venv + install praktika"): + # ci_config was frozen into run metadata by the orchestrator and rides + # on the task, so the job runner installs the pinned Praktika from run + # metadata, not SSM (see ci-config.md). + base_venv, venv_dir = _resolve_runtime( + clone_dir, log, ci_config=task.get("ci_config") or {} + ) if cm_heartbeat is not None: cm_heartbeat.update(phase="writing_task") @@ -864,6 +923,36 @@ def poll(): payload = json.loads(msg["Body"]) log.info("RECEIVED: %s", json.dumps(payload)) + # Read ci_config ONCE for this message: the same snapshot drives the + # controller self-update decision here and (for a workflow) the value + # handle_workflow freezes into run metadata — never re-read SSM, or a + # mid-message change could split the two. + ci_config = _ci_config_for_message(role, payload, log) + + # Idle boundary: converge to the pinned controller version BEFORE doing + # any work for this message. Run the (possibly slow: network download + + # pip) update INSIDE a VisibilityHeartbeat so a long install doesn't let + # the message become visible and get picked up concurrently. If a + # reinstall happened, release the message un-processed (after the + # heartbeat has stopped, so it isn't re-extended) and exit so the + # Restart=always systemd unit relaunches into the new controller. No run + # is interrupted mid-flight. See self_update / ci-config.md. + pin = _controller_version_pin(ci_config) + if pin: + with VisibilityHeartbeat(sqs, queue_url, receipt, visibility): + did_self_update = maybe_self_update(pin, log) + if did_self_update: + try: + sqs.change_message_visibility( + QueueUrl=queue_url, ReceiptHandle=receipt, VisibilityTimeout=0 + ) + except Exception: + log.exception( + "Failed to release message before self-update restart" + ) + log.info("Controller self-update complete; exiting to restart") + return + with VisibilityHeartbeat(sqs, queue_url, receipt, visibility): cleanup_error = _prepare_runner_for_task(role, log) if cleanup_error: @@ -885,7 +974,11 @@ def poll(): if role == ROLE_WORKFLOW: result = handle_workflow( - payload, log, queue_name, receive_count=receive_count + payload, + log, + queue_name, + receive_count=receive_count, + ci_config=ci_config, ) else: result = handle_task( diff --git a/bootstrap/src/praktika_controller/merge.py b/bootstrap/src/praktika_controller/merge.py index 6b30ddd8..641c61e1 100644 --- a/bootstrap/src/praktika_controller/merge.py +++ b/bootstrap/src/praktika_controller/merge.py @@ -35,7 +35,9 @@ import time from pathlib import Path -from praktika_controller.common import load_ci_config +from boto3.s3.transfer import TransferConfig + +from praktika_controller.common import load_ci_config, profile_step # Deterministic identity/dates so the merge sha depends only on the two parents # and the resulting tree, not on wall-clock or runner identity. Must stay @@ -143,6 +145,14 @@ def _apply(path: Path) -> None: # from the merged tree (which carries the target branch's own True settings), # so the rest of the pipeline stays consistent. # + # It also disables the sticky merge base: with force_merge_commit the checkout + # (which may be an upstream-sync head lacking the project's settings override) + # cannot be trusted to resolve the artifact bucket, so the bucket is read from + # the MERGED tree instead (see prepare_repo_snapshot). Sticky-base resolution + # would need that bucket BEFORE the merge, which we no longer have — so it is + # turned off in this mode (accepted limitation; every run just merges into the + # live target tip). + # # ``ci_config`` may be passed pre-resolved by the caller (the controller reads # it once per run and also freezes it into run metadata); fall back to reading # it here when called standalone. @@ -151,9 +161,11 @@ def _apply(path: Path) -> None: if bool(ci_config.get("force_merge_commit")): values["ENABLE_S3_REPO_SNAPSHOT"] = True values["ENABLE_PR_EPHEMERAL_MERGE_COMMIT"] = True + values["STICKY_MERGE_BASE_HOURS"] = 0.0 log.info( "force_merge_commit set in controller config: forcing " - "ENABLE_S3_REPO_SNAPSHOT and ENABLE_PR_EPHEMERAL_MERGE_COMMIT on" + "ENABLE_S3_REPO_SNAPSHOT and ENABLE_PR_EPHEMERAL_MERGE_COMMIT on, " + "disabling sticky merge base (bucket read from merged tree)" ) return values @@ -317,40 +329,64 @@ def _build_and_publish_snapshot(clone_dir, snapshot_sha, is_pr, artifact_bucket, snap_dir = os.path.join(scratch, "repo_snapshot") archive_path = os.path.join(scratch, "repo_snapshot.tar.zst") try: - subprocess.run(["git", "init", "-q", snap_dir], check=True) - # Tag the commit first so it is advertised to the fetch (an unadvertised - # sha would require uploadpack.allowAnySHA1InWant on the source). - _git(["tag", "-f", _REPO_SNAPSHOT_TAG, snapshot_sha], clone_dir) - try: + with profile_step(log, "snapshot: git init + local fetch (depth 1)"): + subprocess.run(["git", "init", "-q", snap_dir], check=True) + # Tag the commit first so it is advertised to the fetch (an unadvertised + # sha would require uploadpack.allowAnySHA1InWant on the source). + _git(["tag", "-f", _REPO_SNAPSHOT_TAG, snapshot_sha], clone_dir) + try: + subprocess.run( + [ + "git", "-C", snap_dir, "fetch", "--depth=1", "-q", + f"file://{os.path.abspath(clone_dir)}", + f"refs/tags/{_REPO_SNAPSHOT_TAG}", + ], + check=True, + ) + finally: + _git(["tag", "-d", _REPO_SNAPSHOT_TAG], clone_dir, check=False) + + with profile_step(log, "snapshot: checkout worktree"): + # Parallel checkout (checkout.workers=0 => one worker per core) to + # write the large working tree to disk faster. subprocess.run( [ - "git", "-C", snap_dir, "fetch", "--depth=1", "-q", - f"file://{os.path.abspath(clone_dir)}", - f"refs/tags/{_REPO_SNAPSHOT_TAG}", + "git", "-C", snap_dir, "-c", "checkout.workers=0", + "checkout", "-q", "--detach", "FETCH_HEAD", ], check=True, ) - finally: - _git(["tag", "-d", _REPO_SNAPSHOT_TAG], clone_dir, check=False) - subprocess.run( - ["git", "-C", snap_dir, "checkout", "-q", "--detach", "FETCH_HEAD"], - check=True, - ) snap_sha = _git_out(["rev-parse", "HEAD"], snap_dir) if snap_sha != snapshot_sha: raise RuntimeError( f"snapshot HEAD {snap_sha} != expected {snapshot_sha}" ) - subprocess.run( - f"tar -C {shlex.quote(snap_dir)} -cf - . | zstd -c -T0 -q > " - f"{shlex.quote(archive_path)}", - shell=True, - check=True, - executable="/bin/bash", - ) - - content_hash = _sha256_file(archive_path) + # Pack and hash in a single pass: stream tar|zstd to this process and + # tee each chunk into both the archive file and the sha256, so the + # content-addressed key needs no separate full-archive read (the archive + # is large for big trees). pipefail so a failing tar isn't masked by a + # succeeding zstd. + with profile_step(log, "snapshot: pack + hash (tar|zstd)"): + proc = subprocess.Popen( + f"set -o pipefail; tar -C {shlex.quote(snap_dir)} -cf - . | zstd -c -T0 -q", + shell=True, + executable="/bin/bash", + stdout=subprocess.PIPE, + ) + h = hashlib.sha256() + assert proc.stdout is not None + with open(archive_path, "wb") as out: + for chunk in iter(lambda: proc.stdout.read(1024 * 1024), b""): + out.write(chunk) + h.update(chunk) + rc = proc.wait() + if rc != 0: + raise RuntimeError( + f"Failed to pack repo snapshot archive (tar|zstd exited {rc})" + ) + content_hash = h.hexdigest() + log.info("[profile] snapshot: archive size %.1f MiB", os.path.getsize(archive_path) / 1048576) # Trust tier: PRs/ (pull_request, fork-reachable) vs REFs/ (push/trusted). # IAM scopes these so a pr-* pool can read but not write REFs/, and a # trusted pool can neither read nor write PRs/. @@ -363,16 +399,40 @@ def _build_and_publish_snapshot(clone_dir, snapshot_sha, is_pr, artifact_bucket, if _s3_object_exists(s3, bucket, key): log.info("Repo snapshot already present: %s", repo_snapshot_key) else: + # Managed multipart upload: stream the archive from disk in parallel + # parts rather than buffering the whole file in RAM and PUTting it + # single-stream (put_object(Body=f.read())). Decisive for large trees + # (e.g. ClickHouse), where the single-stream PUT dominated snapshot + # publish time. + # + # This intentionally drops the old atomic IfNoneMatch="*" write-once. + # Safe for what must be protected — TRUSTED execution (main/release, + # the REFs/ tier): the IAM trust-tier policy (projects.py + # _UNTRUSTED_DENY_WRITE_TRUSTED_STATEMENT) denies untrusted pr-* pools + # PutObject/DeleteObject/AbortMultipartUpload on REFs/*, so no + # untrusted actor can overwrite a trusted snapshot, write-once or not. + # Within the untrusted PRs/ tier one PR can overwrite another PR's + # object (a cross-PR DoS) — an accepted non-threat; the restore-side + # hash check still fails closed, so no tampered bytes ever execute. + # The benign same-key race is fine too (byte-identical archive). try: - with open(archive_path, "rb") as f: - s3.put_object(Bucket=bucket, Key=key, Body=f.read(), IfNoneMatch="*") + with profile_step(log, "snapshot: multipart upload"): + s3.upload_file( + archive_path, + bucket, + key, + Config=TransferConfig( + multipart_threshold=8 * 1024 * 1024, + multipart_chunksize=16 * 1024 * 1024, + max_concurrency=16, + use_threads=True, + ), + ) log.info("Repo snapshot uploaded: %s", repo_snapshot_key) except Exception as e: # noqa: BLE001 - code = getattr(e, "response", {}).get("Error", {}).get("Code", "") - if str(code) in ("PreconditionFailed", "ConditionalRequestConflict"): - # Lost a write-once race — the object exists, which is success. - log.info("Repo snapshot created concurrently: %s", repo_snapshot_key) - elif _s3_object_exists(s3, bucket, key): + if _s3_object_exists(s3, bucket, key): + # Lost a race with a concurrent producer — the object exists, + # which is success. log.info("Repo snapshot created concurrently: %s", repo_snapshot_key) else: raise RuntimeError( @@ -392,14 +452,6 @@ def _s3_object_exists(s3, bucket: str, key: str) -> bool: return False -def _sha256_file(path) -> str: - h = hashlib.sha256() - with open(path, "rb") as f: - for chunk in iter(lambda: f.read(1024 * 1024), b""): - h.update(chunk) - return h.hexdigest() - - def prepare_repo_snapshot(clone_dir, event, settings, s3, log): """Establish the run's single commit and publish its snapshot. @@ -429,13 +481,22 @@ def prepare_repo_snapshot(clone_dir, event, settings, s3, log): # in merge mode, otherwise the head as-is. snapshot_sha = head_sha + # The artifact bucket is infrastructure config, not PR content. In force-merge + # mode it is re-read from the MERGED tree below (post-merge), so it can never + # depend on a PR head that lacks the project's settings override (e.g. an + # upstream-sync branch resolving the public bucket). Sticky base is off in that + # mode, so nothing needs the bucket before the merge; other modes keep the + # value the caller read from the checkout. + artifact_bucket = settings["S3_ARTIFACT_BUCKET"] + if do_merge: base_branch = event.get("base_ref", "") if not base_branch: raise RuntimeError( "Ephemeral merge requested but no base branch on the event" ) - _ensure_base_history(clone_dir, base_branch, log) + with profile_step(log, "merge: ensure base history (unshallow + fetch base)"): + _ensure_base_history(clone_dir, base_branch, log) live_base_sha = _git_out(["rev-parse", f"origin/{base_branch}"], clone_dir) if not live_base_sha: raise RuntimeError(f"Failed to resolve tip of base branch [{base_branch}]") @@ -447,16 +508,17 @@ def prepare_repo_snapshot(clone_dir, event, settings, s3, log): base_sha = live_base_sha sticky_hours = float(settings.get("STICKY_MERGE_BASE_HOURS") or 0) if sticky_hours > 0: - base_sha = _resolve_sticky_base( - s3, - settings["S3_ARTIFACT_BUCKET"], - pr_number, - base_branch, - live_base_sha, - sticky_hours, - clone_dir, - log, - ) + with profile_step(log, "merge: resolve sticky base"): + base_sha = _resolve_sticky_base( + s3, + artifact_bucket, + pr_number, + base_branch, + live_base_sha, + sticky_hours, + clone_dir, + log, + ) log.info( "Ephemeral merge: base [%s] %s (live tip %s) + head %s", base_branch, base_sha[:12], live_base_sha[:12], head_sha[:12], @@ -466,14 +528,15 @@ def prepare_repo_snapshot(clone_dir, event, settings, s3, log): # PR merges cleanly with the CURRENT target HEAD so a green never hides a # real conflict with live main. Non-destructive (git merge-tree, >= 2.38). if base_sha != live_base_sha: - r = _git( - [ - "merge-tree", "--write-tree", "--name-only", - live_base_sha, head_sha, - ], - clone_dir, - check=False, - ) + with profile_step(log, "merge: verify head merges into live tip (merge-tree)"): + r = _git( + [ + "merge-tree", "--write-tree", "--name-only", + live_base_sha, head_sha, + ], + clone_dir, + check=False, + ) if r.returncode == 1: raise MergeConflict( f"PR head {head_sha[:12]} conflicts with the current " @@ -492,13 +555,34 @@ def prepare_repo_snapshot(clone_dir, event, settings, s3, log): f"{r.stdout}\n{r.stderr}", ) - snapshot_sha = _merge_head_into_base( - clone_dir, base_branch, base_sha, head_sha, log - ) + with profile_step(log, "merge: merge head into base"): + snapshot_sha = _merge_head_into_base( + clone_dir, base_branch, base_sha, head_sha, log + ) + + # clone_dir HEAD is now the merge commit, so its ci/settings resolves the + # MERGED tree's artifact bucket (base + head). Re-read it here — this is the + # trusted, PR-independent target for the snapshot upload, correcting a head + # that resolved a different bucket (e.g. an upstream-sync branch missing the + # project's private override). ci_config={} skips the SSM re-read; only the + # bucket is consumed. Falls back to the pre-merge value if the merged tree + # yields nothing. + with profile_step(log, "merge: re-read artifact bucket from merged tree"): + merged_bucket = str( + read_repo_settings(clone_dir, log, ci_config={}).get("S3_ARTIFACT_BUCKET") + or "" + ).strip() + if merged_bucket and merged_bucket != artifact_bucket: + log.info( + "Artifact bucket from merged tree [%s] overrides pre-merge value [%s]", + merged_bucket, artifact_bucket, + ) + artifact_bucket = merged_bucket else: log.info("Repo snapshot: head %s (no merge)", head_sha[:12]) - repo_snapshot_key = _build_and_publish_snapshot( - clone_dir, snapshot_sha, is_pr, settings["S3_ARTIFACT_BUCKET"], s3, log - ) + with profile_step(log, "snapshot: build + publish (total)"): + repo_snapshot_key = _build_and_publish_snapshot( + clone_dir, snapshot_sha, is_pr, artifact_bucket, s3, log + ) return base_sha, snapshot_sha, repo_snapshot_key diff --git a/bootstrap/src/praktika_controller/self_update.py b/bootstrap/src/praktika_controller/self_update.py new file mode 100644 index 00000000..e87dc20b --- /dev/null +++ b/bootstrap/src/praktika_controller/self_update.py @@ -0,0 +1,112 @@ +"""Controller self-upgrade (dev mode). + +The ``praktika_controller_version`` key in the ``ci_config`` SSM parameter pins +the controller wheel a run uses. Unlike the baked AMI/boot version, it can be +rolled forward or back from SSM without rebaking an image. The orchestrator reads +it once from SSM and freezes it into run metadata; every controller (orchestrator +and job runners) converges to the pinned version by reinstalling **in place** into +the system Python and re-launching (the ``Restart=always`` systemd unit relaunches +into the new code). See ``praktika/docs/ci-config.md``. + +The source is any pip install target: a version spec (``praktika-controller==0.1.9``), +a wheel URL (``https://…whl``), or an absolute host path. + +This is a dev-mode mechanism, kept intentionally simple: + +- **Persistence** — the last-installed source string is persisted, so a controller + only reinstalls when the pin actually changes (not on every message). +- **Fail hard** — a bad/unreachable pin is NOT validated, recovered, or rolled + back; pip's error propagates so the failure is loud rather than silently running + the old controller. +""" + +from __future__ import annotations + +import json +import os +import subprocess +import sys +from pathlib import Path + +# Where the last-installed source is persisted. Overridable for tests. +# ``/var/lib/praktika`` is root-writable (the controller runs as root) and already +# used by the system-log streamer. +STATE_PATH = Path( + os.environ.get( + "PRAKTIKA_CONTROLLER_STATE_PATH", + "/var/lib/praktika/controller_version.json", + ) +) + + +def _load_state() -> dict: + try: + data = json.loads(STATE_PATH.read_text(encoding="utf-8")) + if isinstance(data, dict): + return data + except Exception: # noqa: BLE001 - missing/corrupt state -> start fresh + pass + return {"installed_source": ""} + + +def _save_state(state: dict, log=None) -> None: + try: + STATE_PATH.parent.mkdir(parents=True, exist_ok=True) + STATE_PATH.write_text(json.dumps(state), encoding="utf-8") + except Exception as e: # noqa: BLE001 - persistence is best-effort + if log is not None: + log.warning("Could not persist controller version state: %s", e) + + +def _pip_install(source: str, *, python: str, run) -> None: + """Reinstall the controller from ``source`` into ``python``'s site. + + Mirrors the boot-time reinstall: a plain ``--force-reinstall``, retried with + ``--break-system-packages`` when the interpreter is PEP-668 externally-managed + (Ubuntu images). Raises ``CalledProcessError`` if the install fails — a bad pin + fails hard rather than being swallowed.""" + base = [python, "-m", "pip", "install", "--force-reinstall", source] + proc = run(base, capture_output=True, text=True) + if proc.returncode == 0: + return + stderr = proc.stderr or "" + if "externally-managed-environment" in stderr or "--break-system-packages" in stderr: + proc = run(base + ["--break-system-packages"], capture_output=True, text=True) + if proc.returncode == 0: + return + stderr = proc.stderr or "" + raise subprocess.CalledProcessError(proc.returncode, base, stderr=stderr) + + +def maybe_self_update( + desired_source: str, + log, + *, + python: str | None = None, + run=subprocess.run, +) -> bool: + """Converge the controller to ``desired_source`` (dev mode: no validation, no + rollback). + + Returns ``True`` when it reinstalled and the process should restart into the + new code (the caller releases any in-flight message and exits so systemd + relaunches); ``False`` when there is nothing to do (no pin, or already at the + pinned source). A failed install raises — the pin fails hard.""" + desired_source = (desired_source or "").strip() + if not desired_source: + return False # no pin: run whatever is baked/booted + + state = _load_state() + if state.get("installed_source") == desired_source: + # Already at this pin; don't reinstall unless it changes in SSM. Logged so a + # matching pin is visibly honored rather than looking like a no-op. + log.info("Controller already at pinned %s; no reinstall", desired_source) + return False + + python = python or sys.executable or "python3.12" + log.info("Self-updating controller -> %s", desired_source) + _pip_install(desired_source, python=python, run=run) # fail hard on a bad pin + state["installed_source"] = desired_source + _save_state(state, log) + log.info("Controller pinned to %s installed; restarting", desired_source) + return True diff --git a/bootstrap/src/praktika_controller/venv_manager.py b/bootstrap/src/praktika_controller/venv_manager.py index 0b249243..e95aa178 100644 --- a/bootstrap/src/praktika_controller/venv_manager.py +++ b/bootstrap/src/praktika_controller/venv_manager.py @@ -2,6 +2,7 @@ import contextlib import fcntl +import hashlib import os import shutil import subprocess @@ -32,14 +33,40 @@ def ensure_praktika_venv( python_path = str(python_executable or sys.executable) py_tag = f"py{sys.version_info.major}.{sys.version_info.minor}" - env_name = f"praktika-{py_tag}" + # Key the cache by the source too, so a changed source (e.g. a new + # praktika_version pin) resolves to a different venv and is (re)installed, + # instead of silently reusing an existing praktika-* venv built from the old + # source. + env_name = f"praktika-{py_tag}-{_source_key(source)}" venv_dir = cache_root / env_name lock_path = cache_root / f"{env_name}.lock" with _file_lock(lock_path): if _venv_has_praktika(venv_dir): + if is_passthrough(source): + # Immutable source (URL / version spec): the string identity IS the + # content identity, so an existing venv is safe to reuse. + if log is not None: + log.info("Using Praktika from venv %s", venv_dir) + return venv_dir + # Mutable local-path source (e.g. "." / a checkout): the SAME path can + # hold different content across runs (a different restored snapshot), + # so the source string can't prove the cached venv is current. + # Reinstall from source on every task — matching the base-venv overlay + # (_install_runtime_over_base_venv) — so a run never executes stale + # Praktika. --no-deps keeps the venv's baked deps (add a new runtime + # dependency => rebuild the venv). if log is not None: - log.info("Using Praktika from venv %s", venv_dir) + log.info("Reinstalling Praktika from %s into %s", source, venv_dir) + subprocess.run( + _pip_install_cmd( + venv_dir / "bin" / "python", + "--force-reinstall", + "--no-deps", + source, + ), + check=True, + ) return venv_dir if log is not None: @@ -111,9 +138,39 @@ def venv_env( return env +# Requirement-specifier operators (PEP 508). A source containing one of these is +# a version spec (e.g. ``praktika==0.1.9``), not a filesystem path. +_REQUIREMENT_OPERATORS = ("==", ">=", "<=", "~=", "!=", "<", ">") + + +def is_passthrough(source: str) -> bool: + """True when ``source`` is a pip install target that must be passed to pip + verbatim rather than resolved as a local filesystem path: a URL (any + ``scheme://…`` — e.g. an ``https://…whl`` wheel) or a requirement spec (a + package name with a version operator, e.g. ``praktika==0.1.9``). + + Bare package names without an operator are intentionally NOT passthrough: + they are ambiguous with a relative path and version pins always carry an + exact spec, so such inputs stay on the local-path branch.""" + if "://" in source: + return True + return any(op in source for op in _REQUIREMENT_OPERATORS) + + +def _source_key(source: str) -> str: + """Short stable slug of a (normalized) install source, for use in a venv + directory name. A hash keeps arbitrary URLs / long paths bounded and + filesystem-safe while staying 1:1 with the source string.""" + return hashlib.sha256(source.encode("utf-8")).hexdigest()[:12] + + def _normalize_source(source: str) -> str: - # Runtime sources are filesystem paths (a checkout, typically "."); resolve - # to an absolute path for a stable pip target. + # A URL or requirement spec is a pip target in its own right (version pins: + # a wheel URL or ``name==version``); hand it to pip verbatim. Everything else + # is a filesystem path (a checkout, typically "."); resolve to an absolute + # path for a stable pip target. + if is_passthrough(source): + return source return str(Path(source).resolve()) diff --git a/ci/infrastructure/projects.py b/ci/infrastructure/projects.py index 65948dd9..0c5f4140 100644 --- a/ci/infrastructure/projects.py +++ b/ci/infrastructure/projects.py @@ -6,7 +6,7 @@ _PRAKTIKA_PACKAGE_BASE_URL = ( "https://praktika-artifacts-eu-north-1.s3.amazonaws.com/packages" ) -_PRAKTIKA_BASE_VERSION = "0.1.9" +_PRAKTIKA_BASE_VERSION = "0.1.13" # The baked AMI venv pins an exact Praktika version so image builds are # reproducible and a version bump forces a fresh AMI (see _image_builders). _PRAKTIKA_BASE_WHL = ( diff --git a/ci/settings/settings.py b/ci/settings/settings.py index 3d89e987..9201d3ad 100644 --- a/ci/settings/settings.py +++ b/ci/settings/settings.py @@ -25,7 +25,8 @@ class RunnerLabels: STICKY_MERGE_BASE_HOURS = 6 AWS_REGION = "eu-north-1" -AWS_PROFILE = "Box" + +AWS_PROFILE = "SandBox" S3_ARTIFACT_BUCKET = "praktika-artifacts-eu-north-1" S3_REPORT_BUCKET = S3_ARTIFACT_BUCKET diff --git a/ci/tests/test_bootstrap_venv_manager.py b/ci/tests/test_bootstrap_venv_manager.py index 17424b48..51916cbc 100644 --- a/ci/tests/test_bootstrap_venv_manager.py +++ b/ci/tests/test_bootstrap_venv_manager.py @@ -27,21 +27,25 @@ def test_praktika_command_uses_safe_python_path_mode(tmp_path): ] -def test_ensure_praktika_venv_reuses_existing_installed_env(tmp_path, monkeypatch): +def test_ensure_praktika_venv_path_source_reinstalls_on_reuse(tmp_path, monkeypatch): + # A local-path source is mutable (same path, different content across runs), so + # the venv is reused but Praktika is REINSTALLED from source every task — never + # silently served stale. source_dir = tmp_path / "praktika-src" source_dir.mkdir() (source_dir / "setup.py").write_text("from setuptools import setup\n", encoding="utf-8") cache_root = tmp_path / "venvs" + env_name = f"praktika-{_PY_TAG}-{venv_manager._source_key(str(source_dir.resolve()))}" calls = [] def fake_run(cmd, check=False, capture_output=False, text=False, **kwargs): calls.append(cmd) - expected_python = str(cache_root / f"praktika-{_PY_TAG}" / "bin" / "python") + expected_python = str(cache_root / env_name / "bin" / "python") if cmd == [expected_python, "-c", "import praktika"]: return _CompletedProcess( returncode=0 - if (cache_root / f"praktika-{_PY_TAG}" / "bin" / "python").exists() + if (cache_root / env_name / "bin" / "python").exists() else 1 ) if len(cmd) >= 3 and cmd[1:3] == ["-m", "venv"]: @@ -54,25 +58,62 @@ def fake_run(cmd, check=False, capture_output=False, text=False, **kwargs): monkeypatch.setattr(venv_manager.subprocess, "run", fake_run) first = venv_manager.ensure_praktika_venv( - str(source_dir), - cache_root=cache_root, - python_executable="/usr/bin/python3.12", + str(source_dir), cache_root=cache_root, python_executable="/usr/bin/python3.12" ) first_calls = list(calls) - second = venv_manager.ensure_praktika_venv( - str(source_dir), - cache_root=cache_root, - python_executable="/usr/bin/python3.12", + str(source_dir), cache_root=cache_root, python_executable="/usr/bin/python3.12" ) - assert first == second - assert first == (cache_root / f"praktika-{_PY_TAG}").resolve() - - venv_creates = [cmd for cmd in first_calls if len(cmd) >= 3 and cmd[1:3] == ["-m", "venv"]] + assert first == second == (cache_root / env_name).resolve() + # Built once (venv created only on the first call). + venv_creates = [c for c in calls if len(c) >= 3 and c[1:3] == ["-m", "venv"]] assert len(venv_creates) == 1 - assert calls == first_calls + [ - [str(first / "bin" / "python"), "-c", "import praktika"], + # The reuse call force-reinstalls Praktika from the path source. + reuse_calls = calls[len(first_calls):] + reinstall = [ + c + for c in reuse_calls + if c[1:4] == ["-m", "pip", "install"] and "--force-reinstall" in c + ] + assert reinstall, reuse_calls + assert str(source_dir.resolve()) in reinstall[0] + + +def test_ensure_praktika_venv_passthrough_source_reused_without_reinstall( + tmp_path, monkeypatch +): + # A URL / version-spec source is immutable, so an existing venv is reused as-is + # (no reinstall) — the string identity is the content identity. + cache_root = tmp_path / "venvs" + url = "https://x/praktika-0.1.13-py3-none-any.whl" + env_name = f"praktika-{_PY_TAG}-{venv_manager._source_key(url)}" + calls = [] + + def fake_run(cmd, check=False, capture_output=False, text=False, **kwargs): + calls.append(cmd) + expected_python = str(cache_root / env_name / "bin" / "python") + if cmd == [expected_python, "-c", "import praktika"]: + return _CompletedProcess( + returncode=0 if (cache_root / env_name / "bin" / "python").exists() else 1 + ) + if len(cmd) >= 3 and cmd[1:3] == ["-m", "venv"]: + venv_path = Path(cmd[3]) + (venv_path / "bin").mkdir(parents=True, exist_ok=True) + (venv_path / "bin" / "python").write_text("", encoding="utf-8") + return _CompletedProcess() + return _CompletedProcess() + + monkeypatch.setattr(venv_manager.subprocess, "run", fake_run) + + venv_manager.ensure_praktika_venv(url, cache_root=cache_root) + first_calls = list(calls) + venv_manager.ensure_praktika_venv(url, cache_root=cache_root) + + # Reuse call is just the import check — no reinstall for an immutable source. + reuse_calls = calls[len(first_calls):] + assert reuse_calls == [ + [str((cache_root / env_name).resolve() / "bin" / "python"), "-c", "import praktika"] ] def test_ensure_praktika_runtime_uses_base_venv_when_no_source( @@ -249,3 +290,72 @@ def fake_run(cmd, check=False, capture_output=False, text=False, **kwargs): assert "--force-reinstall" in cmd assert "--no-deps" in cmd assert str(source_dir.resolve()) in cmd + + +def test_is_passthrough_url_and_spec_but_not_path(): + assert venv_manager.is_passthrough("https://x/praktika-0.1.9-py3-none-any.whl") + assert venv_manager.is_passthrough("praktika==0.1.9") + assert venv_manager.is_passthrough("praktika>=0.1") + # Bare name and paths are NOT passthrough (ambiguous with a relative path). + assert not venv_manager.is_passthrough("praktika") + assert not venv_manager.is_passthrough(".") + assert not venv_manager.is_passthrough("/opt/praktika/src") + + +def test_normalize_source_passes_through_url_and_spec_resolves_path(tmp_path): + url = "https://x/praktika-0.1.9-py3-none-any.whl" + spec = "praktika==0.1.9" + assert venv_manager._normalize_source(url) == url + assert venv_manager._normalize_source(spec) == spec + # A local path is still resolved to an absolute path. + assert venv_manager._normalize_source(str(tmp_path)) == str(tmp_path.resolve()) + + +def test_source_key_differs_by_source(): + a = venv_manager._source_key("praktika==0.1.9") + b = venv_manager._source_key("praktika==0.2.0") + assert a != b + # Stable for the same source. + assert a == venv_manager._source_key("praktika==0.1.9") + + +def test_ensure_praktika_venv_changed_source_uses_new_venv(tmp_path, monkeypatch): + cache_root = tmp_path / "venvs" + + def fake_run(cmd, check=False, capture_output=False, text=False, **kwargs): + # Report "no praktika yet" so each distinct source triggers a build. + if len(cmd) >= 2 and cmd[-2:] == ["-c", "import praktika"]: + return _CompletedProcess(returncode=1) + if len(cmd) >= 3 and cmd[1:3] == ["-m", "venv"]: + venv_path = Path(cmd[3]) + (venv_path / "bin").mkdir(parents=True, exist_ok=True) + (venv_path / "bin" / "python").write_text("", encoding="utf-8") + return _CompletedProcess() + + monkeypatch.setattr(venv_manager.subprocess, "run", fake_run) + + first = venv_manager.ensure_praktika_venv("praktika==0.1.9", cache_root=cache_root) + second = venv_manager.ensure_praktika_venv("praktika==0.2.0", cache_root=cache_root) + # A different pin resolves to a different venv (so it is actually reinstalled), + # instead of silently reusing the first one. + assert first != second + + +def test_build_venv_installs_url_verbatim(tmp_path, monkeypatch): + cache_root = tmp_path / "venvs" + calls = [] + + def fake_run(cmd, check=False, capture_output=False, text=False, **kwargs): + calls.append(cmd) + # Make the built venv look like it has praktika so ensure_* returns. + return _CompletedProcess() + + monkeypatch.setattr(venv_manager.subprocess, "run", fake_run) + + url = "https://x/praktika-0.1.9-py3-none-any.whl" + venv_manager.ensure_praktika_venv(url, cache_root=cache_root) + install_calls = [ + c for c in calls if len(c) >= 4 and c[1:4] == ["-m", "pip", "install"] + ] + # The wheel URL reaches pip verbatim (not resolved to a filesystem path). + assert any(url in c for c in install_calls) diff --git a/ci/tests/test_controller_job_heartbeat.py b/ci/tests/test_controller_job_heartbeat.py index 072fa12b..5180aed9 100644 --- a/ci/tests/test_controller_job_heartbeat.py +++ b/ci/tests/test_controller_job_heartbeat.py @@ -59,7 +59,7 @@ def test_job_heartbeat_starts_before_runner_setup(monkeypatch, tmp_path): monkeypatch.setattr(controller, "Heartbeat", _FakeHeartbeat) monkeypatch.setattr(controller, "CancelWatchdog", _FakeCancelWatchdog) monkeypatch.setattr(controller, "get_github_token", lambda _region: "token") - monkeypatch.setattr(controller, "_resolve_runtime", lambda *_: ("base", "/venv")) + monkeypatch.setattr(controller, "_resolve_runtime", lambda *_a, **_k: ("base", "/venv")) monkeypatch.setattr(controller, "praktika_command", lambda *_: ["praktika"]) monkeypatch.setattr(controller, "_praktika_env", lambda *_: {}) @@ -133,7 +133,7 @@ def test_cancelled_always_run_job_still_executes(monkeypatch, tmp_path): monkeypatch.setattr(controller, "Heartbeat", _FakeHeartbeat) monkeypatch.setattr(controller, "CancelWatchdog", _FakeCancelWatchdog) monkeypatch.setattr(controller, "get_github_token", lambda _region: "token") - monkeypatch.setattr(controller, "_resolve_runtime", lambda *_: ("base", "/venv")) + monkeypatch.setattr(controller, "_resolve_runtime", lambda *_a, **_k: ("base", "/venv")) monkeypatch.setattr(controller, "praktika_command", lambda *_: ["praktika"]) monkeypatch.setattr(controller, "_praktika_env", lambda *_: {}) diff --git a/ci/tests/test_controller_merge.py b/ci/tests/test_controller_merge.py index 5af61611..3a8e9eca 100644 --- a/ci/tests/test_controller_merge.py +++ b/ci/tests/test_controller_merge.py @@ -44,6 +44,10 @@ def put_object(self, Bucket, Key, Body, IfNoneMatch=None, ContentType=None): raise err self.objects[(Bucket, Key)] = Body + def upload_file(self, Filename, Bucket, Key, Config=None): + with open(Filename, "rb") as f: + self.objects[(Bucket, Key)] = f.read() + def get_object(self, Bucket, Key): if (Bucket, Key) not in self.objects: raise RuntimeError("NoSuchKey") @@ -229,6 +233,46 @@ def test_pr_without_merge_flag_is_plain_head(tmp_path): assert key.startswith("mybucket/prefix/repo-snapshots/v1/PRs/") # still a PR event +def test_merged_tree_bucket_overrides_head(tmp_path): + # Upstream-sync scenario: base (main) carries the project's private bucket via + # an overrides file; the PR head does not (it predates / omits the override). The + # controller must publish to the MERGED tree's bucket, never the head's. + origin = tmp_path / "origin" + origin.mkdir() + _git(origin, "init", "-q", "-b", "main") + os.makedirs(origin / "ci" / "settings") + _write(origin / "ci" / "settings", "settings.py", 'S3_ARTIFACT_BUCKET = "public-bucket"\n') + _git(origin, "add", ".") + _git(origin, "commit", "-q", "-m", "B0") + b0 = _git(origin, "rev-parse", "HEAD") + + # PR head branches from B0 and leaves ci/settings alone -> still resolves public. + _git(origin, "checkout", "-q", "-b", "pr", b0) + _write(origin, "pr.txt", "pr\n") + _git(origin, "add", ".") + _git(origin, "commit", "-q", "-m", "H") + head_sha = _git(origin, "rev-parse", "HEAD") + + # main advances: the private overrides file (resolving the private bucket) lands. + _git(origin, "checkout", "-q", "main") + _write(origin / "ci" / "settings", "private_settings_overrides.py", + 'S3_ARTIFACT_BUCKET = "private-bucket"\n') + _git(origin, "add", ".") + _git(origin, "commit", "-q", "-m", "B1") + + clone = _make_clone(tmp_path, origin, head_sha) + s3 = _FakeS3() + event = {"type": "pull_request", "pr_number": 7, "base_ref": "main"} + # Pre-merge settings resolve the head's public bucket; sticky off (force-merge mode). + settings = dict(PR_MERGE_SETTINGS, S3_ARTIFACT_BUCKET="public-bucket", STICKY_MERGE_BASE_HOURS=0) + + _, _, key = merge.prepare_repo_snapshot(str(clone), event, settings, s3, _Log()) + + # Snapshot published to the merged (private) bucket, and nothing to the public one. + assert key.startswith("private-bucket/repo-snapshots/v1/PRs/") + assert not any(bucket == "public-bucket" for (bucket, _k) in s3.objects) + + def test_read_repo_settings(tmp_path): settings_dir = tmp_path / "ci" / "settings" settings_dir.mkdir(parents=True) @@ -249,6 +293,23 @@ def test_read_repo_settings(tmp_path): assert vals["S3_ARTIFACT_BUCKET"] == "bkt/pfx" +def test_force_merge_commit_disables_sticky(tmp_path): + # force_merge_commit (out-of-repo) forces merge mode ON and, because the bucket + # is then read from the merged tree (post-merge), turns sticky base OFF so + # nothing needs the bucket before the merge. + settings_dir = tmp_path / "ci" / "settings" + settings_dir.mkdir(parents=True) + (settings_dir / "settings.py").write_text( + "STICKY_MERGE_BASE_HOURS = 6\nS3_ARTIFACT_BUCKET = 'bkt'\n" + ) + vals = merge.read_repo_settings( + str(tmp_path), _Log(), ci_config={"force_merge_commit": True} + ) + assert vals["ENABLE_S3_REPO_SNAPSHOT"] is True + assert vals["ENABLE_PR_EPHEMERAL_MERGE_COMMIT"] is True + assert vals["STICKY_MERGE_BASE_HOURS"] == 0.0 # disabled in force-merge mode + + def test_read_repo_settings_defaults_when_absent(tmp_path): # No ci/settings/settings.py -> baked defaults (feature off). vals = merge.read_repo_settings(str(tmp_path), _Log()) diff --git a/ci/tests/test_controller_message_handling.py b/ci/tests/test_controller_message_handling.py index 6d24ef64..05c7df82 100644 --- a/ci/tests/test_controller_message_handling.py +++ b/ci/tests/test_controller_message_handling.py @@ -651,3 +651,89 @@ def test_task_log_capture_records_full_controller_log_to_file(): cap.cleanup() assert not __import__("os").path.exists(cap.path) + + +def _pin(role, payload, log): + return controller._controller_version_pin( + controller._ci_config_for_message(role, payload, log) + ) + + +def test_ci_config_for_message_workflow_reads_ssm_once(monkeypatch): + log = _Log() + reads = [] + monkeypatch.setattr( + controller, + "load_ci_config", + lambda region="", log=None: reads.append(1) + or {"praktika_controller_version": "praktika-controller==5.0"}, + ) + cfg = controller._ci_config_for_message( + controller.ROLE_WORKFLOW, {"type": "pull_request"}, log + ) + assert cfg == {"praktika_controller_version": "praktika-controller==5.0"} + assert controller._controller_version_pin(cfg) == "praktika-controller==5.0" + assert len(reads) == 1 # single SSM read; caller reuses the snapshot + + +def test_ci_config_for_message_rerun_uses_frozen_event_config(monkeypatch): + log = _Log() + # A rerun must NOT re-read SSM; it uses the frozen copy on the event. + monkeypatch.setattr( + controller, + "load_ci_config", + lambda *a, **k: (_ for _ in ()).throw(AssertionError("SSM read on rerun")), + ) + assert ( + _pin( + controller.ROLE_WORKFLOW, + { + "type": "rerun", + "ci_config": {"praktika_controller_version": "praktika-controller==4.0"}, + }, + log, + ) + == "praktika-controller==4.0" + ) + + +def test_ci_config_for_message_runner_reads_task_metadata(monkeypatch): + log = _Log() + monkeypatch.setattr( + controller, + "load_ci_config", + lambda *a, **k: (_ for _ in ()).throw(AssertionError("runner read SSM")), + ) + task = { + "type": "job_task", + "ci_config": {"praktika_controller_version": "praktika-controller==6.0"}, + } + assert _pin(controller.ROLE_RUNNER, task, log) == "praktika-controller==6.0" + # Unset pin -> empty. + assert _pin(controller.ROLE_RUNNER, {"ci_config": {}}, log) == "" + # Explicit empty string is treated as unset (no pin), same as absent. + empty = {"type": "job_task", "ci_config": {"praktika_controller_version": ""}} + assert _pin(controller.ROLE_RUNNER, empty, log) == "" + + +def test_strip_jsonc_preserves_urls_and_strips_comments_and_commas(): + raw = """{ + // pin the controller to an exact wheel + "praktika_controller_version": "https://x.s3.amazonaws.com/praktika_controller-0.1.8.whl", + /* praktika runtime is left floating for now + "praktika_version": "praktika==0.1.9", */ + "force_merge_commit": true, + }""" + data = json.loads(common._strip_jsonc(raw)) + # The // inside the https URL must survive (only out-of-string comments go). + assert data["praktika_controller_version"].startswith("https://") + assert data["praktika_controller_version"].endswith("0.1.8.whl") + # Block-commented field is gone; trailing comma tolerated. + assert "praktika_version" not in data + assert data["force_merge_commit"] is True + + +def test_strip_jsonc_keeps_slashes_and_commas_inside_strings(): + raw = '{"a": "b//c", "list": ["x,]", "y"],}' + data = json.loads(common._strip_jsonc(raw)) + assert data == {"a": "b//c", "list": ["x,]", "y"]} diff --git a/ci/tests/test_controller_self_update.py b/ci/tests/test_controller_self_update.py new file mode 100644 index 00000000..0847fa38 --- /dev/null +++ b/ci/tests/test_controller_self_update.py @@ -0,0 +1,102 @@ +import json +import logging + +import pytest + +from praktika_controller import self_update + + +@pytest.fixture +def state_path(tmp_path, monkeypatch): + path = tmp_path / "controller_version.json" + monkeypatch.setattr(self_update, "STATE_PATH", path) + return path + + +@pytest.fixture +def log(): + return logging.getLogger("test_self_update") + + +class _Proc: + def __init__(self, returncode=0, stdout="", stderr=""): + self.returncode = returncode + self.stdout = stdout + self.stderr = stderr + + +def _read(state_path): + return json.loads(state_path.read_text(encoding="utf-8")) + + +def test_empty_source_is_noop(state_path, log): + calls = [] + assert not self_update.maybe_self_update( + "", log, python="py", run=lambda *a, **k: calls.append(a) or _Proc() + ) + assert calls == [] + assert not state_path.exists() + + +def test_same_source_does_not_reinstall(state_path, log): + state_path.write_text( + json.dumps({"installed_source": "praktika-controller==1.0"}), + encoding="utf-8", + ) + calls = [] + assert not self_update.maybe_self_update( + "praktika-controller==1.0", + log, + python="py", + run=lambda *a, **k: calls.append(a) or _Proc(), + ) + assert calls == [] # persistence: no reinstall unless the pin changes + + +def test_new_source_installs_persists_and_restarts(state_path, log): + installs = [] + + def run(cmd, **kwargs): + installs.append(cmd) + return _Proc(returncode=0) + + assert self_update.maybe_self_update( + "https://x/praktika_controller-3.0-py3-none-any.whl", + log, + python="py", + run=run, + ) + saved = _read(state_path) + assert saved["installed_source"].endswith("praktika_controller-3.0-py3-none-any.whl") + # No post-install verification in dev mode: install then persist + restart. + assert any("install" in c and "--force-reinstall" in c for c in installs) + + +def test_bad_pin_fails_hard(state_path, log): + # A bad/unreachable pin is not validated or rolled back — pip's failure + # propagates so the failure is loud. + def run(cmd, **kwargs): + return _Proc(returncode=1, stderr="ERROR: 404 Not Found") + + with pytest.raises(Exception): + self_update.maybe_self_update( + "https://x/does-not-exist-9.9.whl", log, python="py", run=run + ) + # Nothing persisted as converged. + assert not state_path.exists() + + +def test_break_system_packages_retry_on_pep668(state_path, log): + cmds = [] + + def run(cmd, **kwargs): + cmds.append(cmd) + if "--break-system-packages" not in cmd: + return _Proc(returncode=1, stderr="error: externally-managed-environment") + return _Proc(returncode=0) + + assert self_update.maybe_self_update( + "praktika-controller==3.0", log, python="py", run=run + ) + assert any("--break-system-packages" in c for c in cmds) + assert _read(state_path)["installed_source"] == "praktika-controller==3.0" diff --git a/ci/tests/test_workflow_routing.py b/ci/tests/test_workflow_routing.py index 66a46660..7a808815 100644 --- a/ci/tests/test_workflow_routing.py +++ b/ci/tests/test_workflow_routing.py @@ -1,3 +1,4 @@ +import pytest from praktika import Job, Workflow from praktika.mangle import _update_workflow_with_native_jobs from praktika.orchestrator import find_workflows_for_event @@ -28,7 +29,7 @@ def test_default_orchestrator_skips_base_workflows(monkeypatch): monkeypatch.setenv("PRAKTIKA_CONTROLLER_QUEUE", "workflow-orchestrator") monkeypatch.setattr( "praktika.orchestrator._get_workflows", - lambda: [default_workflow, base_workflow], + lambda *a, **k: [default_workflow, base_workflow], ) matched = find_workflows_for_event(event) @@ -44,7 +45,7 @@ def test_workflow_name_filter_selects_one_matching_workflow(monkeypatch): monkeypatch.setenv("PRAKTIKA_CONTROLLER_QUEUE", "workflow-orchestrator") monkeypatch.setattr( "praktika.orchestrator._get_workflows", - lambda: [pr_fast, pr_full], + lambda *a, **k: [pr_fast, pr_full], ) matched = find_workflows_for_event(event, workflow_name="PR Full") @@ -59,7 +60,7 @@ def test_workflow_name_filter_ignores_orchestrator_pool_filter(monkeypatch): monkeypatch.setenv("PRAKTIKA_CONTROLLER_QUEUE", "workflow-orchestrator") monkeypatch.setattr( "praktika.orchestrator._get_workflows", - lambda: [base_workflow], + lambda *a, **k: [base_workflow], ) matched = find_workflows_for_event(event, workflow_name="Praktika CI") @@ -75,7 +76,7 @@ def test_base_orchestrator_skips_default_workflows(monkeypatch): monkeypatch.setenv("PRAKTIKA_CONTROLLER_QUEUE", "workflow-orchestrator-base") monkeypatch.setattr( "praktika.orchestrator._get_workflows", - lambda: [default_workflow, base_workflow], + lambda *a, **k: [default_workflow, base_workflow], ) matched = find_workflows_for_event(event) @@ -103,7 +104,7 @@ def test_schedule_event_routes_to_ignition_workflow(monkeypatch): event = {"type": "schedule", "head_ref": "main", "workflow_name": "Nightly"} monkeypatch.setenv("PRAKTIKA_CONTROLLER_QUEUE", "workflow-orchestrator") - monkeypatch.setattr("praktika.orchestrator._get_workflows", lambda: [wf]) + monkeypatch.setattr("praktika.orchestrator._get_workflows", lambda *a, **k: [wf]) matched = find_workflows_for_event(event) @@ -115,7 +116,7 @@ def test_dispatch_event_routes_to_ignition_workflow(monkeypatch): event = {"type": "dispatch", "head_ref": "main", "workflow_name": "Release"} monkeypatch.setenv("PRAKTIKA_CONTROLLER_QUEUE", "workflow-orchestrator") - monkeypatch.setattr("praktika.orchestrator._get_workflows", lambda: [wf]) + monkeypatch.setattr("praktika.orchestrator._get_workflows", lambda *a, **k: [wf]) matched = find_workflows_for_event(event) @@ -133,7 +134,7 @@ def test_dispatch_without_branches_runs_on_any_ref(monkeypatch): } monkeypatch.setenv("PRAKTIKA_CONTROLLER_QUEUE", "workflow-orchestrator") - monkeypatch.setattr("praktika.orchestrator._get_workflows", lambda: [wf]) + monkeypatch.setattr("praktika.orchestrator._get_workflows", lambda *a, **k: [wf]) matched = find_workflows_for_event(event) @@ -151,7 +152,7 @@ def test_dispatch_with_branches_restricts_ref(monkeypatch): } monkeypatch.setenv("PRAKTIKA_CONTROLLER_QUEUE", "workflow-orchestrator") - monkeypatch.setattr("praktika.orchestrator._get_workflows", lambda: [wf]) + monkeypatch.setattr("praktika.orchestrator._get_workflows", lambda *a, **k: [wf]) matched = find_workflows_for_event(event) @@ -166,7 +167,7 @@ def test_ignition_event_name_does_not_bypass_pool_routing(monkeypatch): event = {"type": "schedule", "head_ref": "main", "workflow_name": "Nightly"} monkeypatch.setenv("PRAKTIKA_CONTROLLER_QUEUE", "workflow-orchestrator") # default - monkeypatch.setattr("praktika.orchestrator._get_workflows", lambda: [wf]) + monkeypatch.setattr("praktika.orchestrator._get_workflows", lambda *a, **k: [wf]) assert find_workflows_for_event(event) == [] @@ -177,7 +178,7 @@ def test_ignition_event_name_runs_on_its_own_pool(monkeypatch): event = {"type": "schedule", "head_ref": "main", "workflow_name": "Nightly"} monkeypatch.setenv("PRAKTIKA_CONTROLLER_QUEUE", "workflow-orchestrator-base") - monkeypatch.setattr("praktika.orchestrator._get_workflows", lambda: [wf]) + monkeypatch.setattr("praktika.orchestrator._get_workflows", lambda *a, **k: [wf]) assert [w.name for w in find_workflows_for_event(event)] == ["Nightly"] @@ -190,7 +191,7 @@ def test_explicit_name_arg_still_bypasses_pool_routing(monkeypatch): event = {"type": "schedule", "head_ref": "main"} monkeypatch.setenv("PRAKTIKA_CONTROLLER_QUEUE", "workflow-orchestrator") # default - monkeypatch.setattr("praktika.orchestrator._get_workflows", lambda: [wf]) + monkeypatch.setattr("praktika.orchestrator._get_workflows", lambda *a, **k: [wf]) matched = find_workflows_for_event(event, workflow_name="Nightly") assert [w.name for w in matched] == ["Nightly"] @@ -202,7 +203,7 @@ def test_schedule_without_workflow_name_is_skipped(monkeypatch): event = {"type": "schedule", "head_ref": "main"} monkeypatch.setenv("PRAKTIKA_CONTROLLER_QUEUE", "workflow-orchestrator") - monkeypatch.setattr("praktika.orchestrator._get_workflows", lambda: [wf]) + monkeypatch.setattr("praktika.orchestrator._get_workflows", lambda *a, **k: [wf]) matched = find_workflows_for_event(event) @@ -215,7 +216,7 @@ def test_schedule_message_name_filters_to_one_workflow(monkeypatch): event = {"type": "schedule", "head_ref": "main", "workflow_name": "Nightly B"} monkeypatch.setenv("PRAKTIKA_CONTROLLER_QUEUE", "workflow-orchestrator") - monkeypatch.setattr("praktika.orchestrator._get_workflows", lambda: [a, b]) + monkeypatch.setattr("praktika.orchestrator._get_workflows", lambda *args, **kw: [a, b]) matched = find_workflows_for_event(event) @@ -231,7 +232,7 @@ def test_gh_actions_schedule_workflow_is_skipped(monkeypatch): event = {"type": "schedule", "head_ref": "main", "workflow_name": "GH Nightly"} monkeypatch.setenv("PRAKTIKA_CONTROLLER_QUEUE", "workflow-orchestrator") - monkeypatch.setattr("praktika.orchestrator._get_workflows", lambda: [wf]) + monkeypatch.setattr("praktika.orchestrator._get_workflows", lambda *a, **k: [wf]) matched = find_workflows_for_event(event) @@ -261,3 +262,113 @@ def test_native_jobs_can_follow_base_runner_override(): assert workflow.jobs[0].runs_on == ["arm-2xsmall-base"] assert workflow.jobs[-1].name == Settings.FINISH_WORKFLOW_JOB_NAME assert workflow.jobs[-1].runs_on == ["arm-2xsmall-base"] + + +_GOOD_WF_FILE = '''\ +from praktika import Job, Workflow + +WORKFLOWS = [ + Workflow.Config( + name="Good WF", + event=Workflow.Event.PULL_REQUEST, + base_branches=["main"], + jobs=[Job.Config(name="User Job", runs_on=["arm-2xsmall"], command="true")], + ) +] +''' + +# References an attribute that does not exist on Workflow.Engine — mirrors a pool +# whose baked praktika predates a newer engine used by a workflow file. +_BROKEN_WF_FILE = '''\ +from praktika import Workflow + +_ = Workflow.Engine.DOES_NOT_EXIST +WORKFLOWS = [] +''' + + +def _write_workflows_dir(tmp_path): + (tmp_path / "good_wf.py").write_text(_GOOD_WF_FILE, encoding="utf-8") + (tmp_path / "broken_wf.py").write_text(_BROKEN_WF_FILE, encoding="utf-8") + return tmp_path + + +def test_get_workflows_skips_broken_file_and_records_error(tmp_path, monkeypatch): + from praktika import mangle + + _write_workflows_dir(tmp_path) + monkeypatch.setattr(Settings, "WORKFLOWS_DIRECTORY", str(tmp_path)) + + errors = [] + res = mangle._get_workflows(name="Good WF", _load_errors_out=errors) + + # The good workflow still loads despite the broken sibling file. + assert [wf.name for wf in res] == ["Good WF"] + # The broken file is surfaced (filename + error) rather than silently dropped. + assert [f for f, _ in errors] == ["broken_wf.py"] + assert "DOES_NOT_EXIST" in errors[0][1] + + +def test_get_workflows_all_files_broken_returns_empty_with_errors(tmp_path, monkeypatch): + from praktika import mangle + + # Only broken files: must NOT raise "no workflow found" (which the caller can't + # tell from a real misconfig). Return empty + the collected load errors so the + # orchestrator can surface them and finalize the bootstrap check. + (tmp_path / "broken_wf.py").write_text(_BROKEN_WF_FILE, encoding="utf-8") + monkeypatch.setattr(Settings, "WORKFLOWS_DIRECTORY", str(tmp_path)) + + errors = [] + res = mangle._get_workflows(_load_errors_out=errors) + assert res == [] + assert [f for f, _ in errors] == ["broken_wf.py"] + + +def test_get_workflows_reraises_broken_file_during_validation(tmp_path, monkeypatch): + from praktika import mangle + + _write_workflows_dir(tmp_path) + monkeypatch.setattr(Settings, "WORKFLOWS_DIRECTORY", str(tmp_path)) + + # Validation must NOT swallow a broken workflow file (it runs under the + # checkout's own praktika, so an import error is a genuine bug to surface). + with pytest.raises(AttributeError): + mangle._get_workflows(_for_validation_check=True) + + +def test_orchestrate_posts_failure_check_for_workflow_load_error(monkeypatch): + # A skipped (unimportable) workflow file must surface as a real completed + # failure check — regression guard: CheckRun exposes start()+complete(), not a + # create_completed() one-shot. + from praktika import orchestrator + from praktika.orchestrator import check_run + + posted = {} + + class _FakeCheck: + def complete(self, conclusion, output=None, details_url=None): + posted["conclusion"] = conclusion + posted["output"] = output + + def fake_start(cls, token, repo, head_sha, name, details_url=None, with_cancel_action=True): + posted["name"] = name + posted["with_cancel_action"] = with_cancel_action + return _FakeCheck() + + monkeypatch.setattr(check_run.CheckRun, "start", classmethod(fake_start)) + + def fake_find(event, workflow_name=None, _load_errors_out=None): + if isinstance(_load_errors_out, list): + _load_errors_out.append(("ignition_dispatch.py", "AttributeError: GH_IGNITION")) + return [] + + monkeypatch.setattr(orchestrator, "find_workflows_for_event", fake_find) + + event = {"type": "pull_request", "repo": "o/r", "head_sha": "abc123"} + rc = orchestrator._orchestrate_event(event, gh_token="tok", run_id="1", ci=True) + + assert rc == 0 + assert posted["name"] == "Workflow load error: ignition_dispatch.py" + assert posted["conclusion"] == "failure" + assert posted["with_cancel_action"] is False + assert "GH_IGNITION" in posted["output"]["summary"] diff --git a/ci/workflows/praktika_pr_advanced.py b/ci/workflows/praktika_pr_advanced.py index 10293288..54698982 100644 --- a/ci/workflows/praktika_pr_advanced.py +++ b/ci/workflows/praktika_pr_advanced.py @@ -8,7 +8,7 @@ from ci.settings.settings import RunnerLabels from praktika.settings import Settings -_HEAD_PRAKTIKA_VERSION = "0.1.13" +_HEAD_PRAKTIKA_VERSION = "0.1.14" artifact = Artifact.Config(name="greet", type=Artifact.Type.S3, path="./artifact.txt") diff --git a/praktika/__main__.py b/praktika/__main__.py index 3ae4aa29..51a1ac8f 100644 --- a/praktika/__main__.py +++ b/praktika/__main__.py @@ -192,7 +192,8 @@ def create_parser(): help=( "Process only specified components (e.g. html ImageBuilder AMI VPC LaunchTemplate AutoScalingGroup Lambda DedicatedHost EC2Instance). " "With --deploy: deploys only these components or uploads html report. " - "With --destroy-runtime/--destroy-all: deletes only the selected component types." + "With --destroy-runtime/--destroy-all: deletes only the selected component types. " + "With --restart-instances: refreshes only the selected ASGs (e.g. DockerProxy)." ), nargs="+", type=str, @@ -200,7 +201,7 @@ def create_parser(): ) _infra_parser.add_argument( "--restart-instances", - help="Trigger an instance refresh on all ASGs, replacing all EC2 instances with the current launch template version", + help="Trigger an instance refresh on ASGs, replacing their EC2 instances with the current launch template version. Refreshes all ASGs by default; scope with --only (e.g. --only DockerProxy) to roll a single pool without disturbing the others", action="store_true", default=False, ) @@ -400,7 +401,7 @@ def main(argv=None): if args.restart_instances: from .mangle import _get_infra_config - _get_infra_config(project).restart_instances() + _get_infra_config(project).restart_instances(only=args.only) if args.verify: from .mangle import _get_infra_config diff --git a/praktika/docs/ci-config.md b/praktika/docs/ci-config.md index 4c16fddb..8862d2bd 100644 --- a/praktika/docs/ci-config.md +++ b/praktika/docs/ci-config.md @@ -6,7 +6,11 @@ Store, outside the repository. It is the home for CI settings that must apply controls, not the code under test. - **Parameter name:** `{PROJECT_SLUG}-ci-config` (e.g. `myproject-ci-config`). -- **Type:** `String`, holding a JSON **object**. +- **Type:** `String`, holding a JSON **object**. JSONC is tolerated — `//` and + `/* */` comments and trailing commas are stripped before parsing (`_strip_jsonc`), + so you can comment a field out without silently invalidating the whole parameter + (a strict parse error reads as `{}` — every feature off). Comments/commas inside + string values (e.g. a `//` in an `https://` URL) are preserved. - **Region:** the project's region (`AWS_DEFAULT_REGION` / `AWS_REGION`). - **Reader:** `praktika_controller.common.load_ci_config`. @@ -42,17 +46,14 @@ SSM (rather than an instance tag or env var) is used deliberately: ## Intended scope (roadmap) -Today the CI config carries a single migration switch (below). It is meant to -grow into the general surface for out-of-repo CI settings, for example: +Today the CI config carries a migration switch and version pins (below). It is +meant to grow into the general surface for out-of-repo CI settings, for example: - **More CI-wide operational toggles** — rollout flags, feature gates, temporary overrides during incidents or migrations. - **Per-user settings** — a `users` sub-object keyed by GitHub login, letting individual developers opt into experimental behavior for their own PRs without touching the repo (e.g. `{"users": {"alice": {"...": true}}}`). -- **Version pinning** — pin the praktika / praktika-controller version a run uses, - read once and preserved in run metadata (see [Proposed: version - pinning](#proposed-version-pinning-not-yet-implemented)). When adding a key, keep the same contract: optional, safe-by-absence, and documented here. @@ -86,20 +87,27 @@ the ephemeral merge happens, the run reinstalls praktika from the **merged tree* which carries the target branch's own (True) settings — so the config workflow, DAG, and every job stay consistent without any further override. +**Limitation — sticky merge base is disabled.** `force_merge_commit` also forces +`STICKY_MERGE_BASE_HOURS` to `0`, so every run merges the PR head into the *live* +target tip (the sticky-base optimization, which reuses a pinned target commit +across a PR's rapid re-pushes to keep the digest cache warm, is off). This is +deliberate: with `force_merge_commit` the checkout (e.g. an upstream-sync head) +can't be trusted to resolve `S3_ARTIFACT_BUCKET`, so the bucket is read from the +**merged tree** *after* the merge. Sticky-base resolution would need that bucket +*before* the merge (it reads/writes an S3 pin), which is no longer available — so +it is turned off in this mode. See `praktika_controller.merge.read_repo_settings` +and `prepare_repo_snapshot`. + **Lifecycle.** This is a migration lever: set it when you enable the feature, leave it on until old branches age out, then delete the parameter (or set it to `false`). Because it's read per run, both flipping it on and removing it take effect on the next run. -## Proposed: version pinning (not yet implemented) - -> Status: **design only.** Nothing below is wired up yet — the sole implemented -> setting today is `force_merge_commit`. This section records the intended shape -> so the config schema grows coherently. +### `praktika_version` / `praktika_controller_version` (version pins) -Goal: pin the praktika (and praktika-controller) version a run uses, so an -unexpected upgrade cannot break an already-started run, and so the version can be -rolled out from SSM instead of rebaking the AMI. Proposed keys: +Pin the praktika and/or praktika-controller version a run uses, so an unexpected +upgrade cannot break an already-started run and so a version can be rolled out (or +rolled back) from SSM instead of rebaking the AMI: ```json { @@ -108,70 +116,107 @@ rolled out from SSM instead of rebaking the AMI. Proposed keys: } ``` -Each `` supports three interchangeable forms, all of which `pip install` -already accepts as-is: +Both are optional and independent, and an empty string (`""`) is treated exactly +like "not set" (no pin). Each `` supports three interchangeable forms, all +of which `pip install` accepts as-is (detected by `venv_manager.is_passthrough`): -- **version** — a released spec, e.g. `"praktika==0.1.9"` (once published to an - index). +- **version** — a released spec, e.g. `"praktika==0.1.9"` (must include an exact + `==`; a bare name is treated as a path). - **https path** — a wheel URL, e.g. `"https://.../praktika-0.1.9-py3-none-any.whl"`. -- **repo path** — a filesystem path/checkout, e.g. `"."` or `/opt/praktika/src` - (the current `runtime_source` behavior). - -### Two different lifecycles (do not conflate) - -- **`praktika_version` → per-run pin.** Read once at run start, frozen into the - run metadata (the `_Environment` / `task` / `state.json` carrier — same pattern - as `SNAPSHOT_SHA` / `WORKFLOW_START_TIME`), and every job of the run sources its - runtime from that frozen value. This is what protects a started run from a - mid-run upgrade: all jobs agree on one version regardless of what changes in SSM - afterward. In S3-snapshot mode a repo-sourced praktika is *already* content-pinned - by the snapshot; this closes the gap for the URL / version / moving-source forms. -- **`praktika_controller_version` → global desired-state, NOT a per-run pin.** One - controller process serves many runs off the queue, so it cannot run a different - controller version per in-flight run. Instead the controller converges to the - SSM-desired version by self-reinstalling **between runs** (at an idle boundary), - never mid-run. It still protects in-flight runs (the reinstall/re-exec happens - when the controller is idle), but the mechanism is a rolling self-update. - -### Sketch of the mechanism - -- **`venv_manager._normalize_source`** must branch on the form (URL scheme or - requirement spec → pass through verbatim; otherwise resolve as a local path). - This is the only code gap for the three forms; pip handles the rest. -- **Version introspection already exists** (`praktika/version.py`: - `current_praktika_version`, `current_praktika_controller_version`, - `version_key`). Mismatch detection compares the running version to the desired - one. -- **Controller self-reinstall** reads `praktika_controller_version` from SSM - (trusted infra, read before any PR code runs — keep it in SSM, never in - PR-influenced metadata), and on mismatch installs the desired form into a fresh - overlay venv and re-execs into it (mirroring `_install_runtime_over_base_venv`), - rather than mutating the live system-python in place. - -### Risks / guards this needs before shipping - -- **Crash-loop protection** — a bad version otherwise makes every fresh instance - reinstall → crash → restart → reinstall forever. Fall back to the baked version - on install/import/health failure, cap attempts, persist last-known-good. -- **Install atomicity** — install into an overlay venv and re-exec into it; never - half-mutate the running controller's env. -- **Idle boundary** — only reinstall/re-exec with no message in flight (or at - process start, before claiming work). -- **Bootstrapping** — only controllers that already ship this logic can - self-update; the first rollout is still a normal AMI/deploy. - -Because parts differ in risk, the intended rollout order is: (1) `_normalize_source` -three-form support + freeze `praktika_version` into run metadata (low risk), then -(2) controller self-reinstall as a separate, carefully-guarded change. +- **repo path** — a filesystem path/checkout, e.g. `"."` or `/opt/praktika/src`. + For `praktika_version`, relative paths resolve against the run's checkout. For + `praktika_controller_version`, **only absolute paths** are allowed (a URL or spec + otherwise): controller self-update runs *before* any checkout exists, so a + relative path has nothing to resolve against and is rejected. + +Both pins are read **once from SSM by the orchestrator** and frozen into run +metadata via the same `ci_config` carrier as `force_merge_commit` (§ *How the +config reaches jobs*), so every controller of the run obeys the value the +orchestrator pinned — **read from run metadata, not SSM**. + +#### When it applies (praktika development setup) + +Version pinning targets the **praktika development setup**: pools whose runtime is +*not* already pinned on the runner — i.e. the install source points at a **floating +"latest"** (a `.../latest/…whl` alias, a moving branch checkout, or a base venv +built to track head). If a pool already pins an exact version on the runner (a +versioned base venv / wheel), that fixed version is what runs; a `ci_config` pin is +redundant there. + +Its purpose in that floating setup is twofold: + +- **Pin one version across the whole run fleet.** Without a pin, each instance + resolves "latest" independently, so a fleet can end up straddling versions + (some runners on the wheel published a minute ago, others on the previous one) — + and within a single run the orchestrator and its job runners could disagree. + A pin freezes one version into run metadata so the orchestrator and every job + runner of the run use exactly the same one. +- **Change the version without an infra update.** Rolling a floating setup forward + or back otherwise means republishing the "latest" wheel and/or rebaking the AMI. + A pin moves that control to a single SSM edit that the *next* run picks up — no + wheel republish, no AMI rebake, no instance replacement. + +#### `praktika_version` (per-run runtime pin) + +Read once at run start, frozen into run metadata, and every job runner installs +its praktika runtime from that frozen value (via +`controller._resolve_runtime_source`, taking precedence over a pool's +`praktika_runtime_source` tag). All jobs of a run agree on one version regardless +of what changes in SSM afterward. In S3-snapshot mode a repo-sourced praktika is +already content-pinned by the snapshot; this closes the gap for the URL / version +forms. Installed into the base-venv overlay with `--no-deps`, so the base venv must +already carry praktika's runtime dependencies (a new dependency needs an AMI +rebake, as today). + +#### `praktika_controller_version` (per-run controller pin, self-upgrade) + +The controller is the `praktika-controller` wheel installed into the system +`python3.12` and run by systemd (`Restart=always`). When the frozen +`praktika_controller_version` differs from the source the controller last +installed, the controller **self-upgrades**: at the idle boundary (a message is +received but not yet processed), it `pip install --force-reinstall`s the pinned +source into system python, persists the new source, releases the message back to +the queue un-processed, and exits — the `Restart=always` unit relaunches into the +new code, which re-receives the message and proceeds. No run is interrupted +mid-flight. Mechanism lives in `praktika_controller.self_update.maybe_self_update`, +wired into `controller.poll()`. + +This is a **dev-mode** mechanism, kept deliberately simple: + +- **Source of the desired version.** The orchestrator (workflow role) reads it + from SSM (`load_ci_config`) for a fresh run, or the frozen copy on a rerun event. + Job runners read it from the task's frozen `ci_config` — never SSM — so all + runners of a run converge to the version the orchestrator pinned. +- **Persistence.** The last-installed source string is persisted to + `/var/lib/praktika/controller_version.json` (override with + `PRAKTIKA_CONTROLLER_STATE_PATH`), so the controller reinstalls only when the pin + string changes, not on every message. +- **Fail hard.** The pin is not validated, version-checked, or rolled back. A + bad/unreachable source makes pip fail and the error propagates (the message is + retried / the instance replaced by the normal infra-failure path) rather than + silently running the old controller. Use an exact, reachable pin. +- **Bootstrapping.** Only controllers that already ship this logic can self-update; + the first rollout is a normal AMI / boot-time wheel install, which also remains + the fallback when no pin is set. + +> **Caveat — runner flapping.** Because the controller version is frozen per run, a +> runner serving tasks from two concurrently-active runs pinned to *different* +> controller versions will reinstall/restart back and forth. In practice SSM is +> stable and all live runs share one value; changing SSM only affects new runs +> while in-flight runs keep their frozen value. +> +> **Caveat — moving sources.** Persistence is keyed on the source *string*, so a +> mutable `…/latest/…whl` URL is not detected as "changed". Pin to an immutable +> exact version or versioned URL. ## Setting / clearing the parameter ```bash -# Enable (migration on) +# Enable (migration on) / pin versions aws ssm put-parameter \ --name "{PROJECT_SLUG}-ci-config" \ --type String \ - --value '{"force_merge_commit": true}' \ + --value '{"force_merge_commit": true, "praktika_version": "praktika==0.1.9", "praktika_controller_version": "praktika-controller==0.1.8"}' \ --overwrite \ --region "$AWS_REGION" diff --git a/praktika/infrastructure/cloud.py b/praktika/infrastructure/cloud.py index 644d6889..107cb0ef 100644 --- a/praktika/infrastructure/cloud.py +++ b/praktika/infrastructure/cloud.py @@ -40,6 +40,28 @@ from .sqs_queue import SQSQueue +# Seed value for the {slug}-ci-config SSM parameter created at deploy time: an +# effectively-empty config (parses to {}) documenting every supported key, +# commented out. JSONC is tolerated by the reader (praktika_controller.common. +# load_ci_config), so operators uncomment a line to enable it. Keep descriptions +# short; the full contract lives in praktika/docs/ci-config.md. +_CI_CONFIG_TEMPLATE = """\ +{ + // Praktika out-of-repo CI config. JSONC ok (// comments, trailing commas). + // Uncomment a key to enable it. Docs: praktika/docs/ci-config.md + + // Force ephemeral PR merge + repo snapshot on every run (migration lever). + // "force_merge_commit": true, + + // Pin praktika runtime for a run: version spec / wheel URL / repo path. + // "praktika_version": "praktika==0.1.14", + + // Pin controller wheel: version spec / wheel URL / absolute host path. + // "praktika_controller_version": "praktika-controller==0.1.8" +} +""" + + class CloudInfrastructure: SLACK_APP_LAMBDAS = [lambda_app_config, lambda_worker_config] @@ -1738,8 +1760,40 @@ def _wants(type_name: str, *aliases: str) -> bool: asg_config.deploy() deployed_asg_configs.append(asg_config) + # Seed the out-of-repo CI config SSM parameter (create-if-absent, so an + # operator's existing value is never clobbered by a redeploy). + if _wants("CIConfig", "ci-config", "ciconfig"): + print("\n" + "=" * 60) + print("Deploying CI config SSM parameter") + print("=" * 60) + self._ensure_ci_config_parameter() + self._print_deployment_warnings(deployed_asg_configs) + def _ensure_ci_config_parameter(self): + """Create the ``{slug}-ci-config`` SSM parameter with a commented-out + template (parses to ``{}`` — every feature off) if it does not exist. + + Create-if-absent only: a redeploy must never overwrite an operator's + live value (it's edited out-of-band, not from the repo). The name + matches what the controller reads (``PRAKTIKA_PROJECT_SLUG`` = the + project prefix). See praktika/docs/ci-config.md.""" + region = self._settings.AWS_REGION + name = f"{self._project_prefix()}-ci-config" + client = aws_client("ssm", region, f"{self._project_prefix()}-ci-config") + try: + client.put_parameter( + Name=name, + Type="String", + Value=_CI_CONFIG_TEMPLATE, + Description=( + "Praktika out-of-repo CI config; see praktika/docs/ci-config.md" + ), + ) + print(f"Created SSM parameter {name} (commented-out template)") + except client.exceptions.ParameterAlreadyExists: + print(f"SSM parameter {name} already exists; leaving its value unchanged") + def destroy_runtime( self, force: bool = True, @@ -2640,13 +2694,50 @@ def _delete_bucket(bucket_name=name): ) print("=" * 60) - def restart_instances(self): - """Trigger an instance refresh on all configured ASGs.""" + def restart_instances(self, only: Optional[List[str]] = None): + """Trigger an instance refresh on the configured ASGs. + + With no ``only`` filter this refreshes every ASG (all runner and + orchestrator pools too), which rolls their instances and kills any + in-flight jobs -- use it deliberately. Pass ``only`` to scope the + refresh, e.g. ``--only DockerProxy`` to roll just the DockerHub proxy + (a stateless single-instance cache, safe to replace) onto its current + launch template version. + """ self._verify_account() - if not self.autoscaling_groups: - print("No ASGs configured") + + only_set = { + s.strip().lower() + for s in (only or []) + if isinstance(s, str) and s.strip() + } + + def _wants(type_name: str, *aliases: str) -> bool: + if not only_set: + return True + keys = {type_name.lower(), *{a.lower() for a in aliases if a}} + return bool(keys & only_set) + + # De-dupe by ASG name: the DockerProxy ASG is also present in + # self.autoscaling_groups, so a plain "all" pass would list it twice. + selected: Dict[str, "AutoScalingGroup.Config"] = {} + + if self.docker_proxy and _wants( + "DockerProxy", "docker-proxy", "dockerproxy", "dockerhub-proxy" + ): + selected[self.docker_proxy.autoscaling_group.name] = ( + self.docker_proxy.autoscaling_group + ) + + if _wants("AutoScalingGroup", "AutoScalingGroups", "ASG", "ASGs"): + for asg_config in self.autoscaling_groups: + selected.setdefault(asg_config.name, asg_config) + + if not selected: + print(f"No ASGs match the selection: {sorted(only_set)}") return - for asg_config in self.autoscaling_groups: + + for asg_config in selected.values(): asg_config.region = self._settings.AWS_REGION print("\n" + "=" * 60) print(f"Restarting instances in ASG: {asg_config.name}") diff --git a/praktika/infrastructure/launch_template.py b/praktika/infrastructure/launch_template.py index 965d3bfe..4ae0720a 100644 --- a/praktika/infrastructure/launch_template.py +++ b/praktika/infrastructure/launch_template.py @@ -1,4 +1,5 @@ import base64 +import re from dataclasses import dataclass, field from typing import TYPE_CHECKING, Any, Dict, List, Optional @@ -176,10 +177,13 @@ def _resolve_ready_ami_from_pipeline(client, pipeline_arn: str, label: str) -> s ) return self.image_id - # Detect architecture from instance type: Graviton families end in 'g' - # (t4g, m6g, c6g, r6g, ...). Everything else is x86_64. + # Detect architecture from instance type: Graviton (arm64) families + # have a 'g' immediately after the generation number, optionally + # followed by suffix letters (t4g, c6g, c7gn, c7gd, m7gd, im4gn, ...). + # Everything else is x86_64. Matching just endswith("g") missed the + # 'gn'/'gd' network/disk-optimized variants. family = (self.instance_type or "").split(".")[0] - is_arm = family.endswith("g") + is_arm = bool(re.search(r"[0-9]g[a-z]*$", family)) if is_arm: from .native.configs import resolve_al2023_arm64_ami self.image_id = resolve_al2023_arm64_ami(self.region) diff --git a/praktika/infrastructure/native/docker_proxy.md b/praktika/infrastructure/native/docker_proxy.md index 33eb05ae..bf6c40bf 100644 --- a/praktika/infrastructure/native/docker_proxy.md +++ b/praktika/infrastructure/native/docker_proxy.md @@ -6,11 +6,16 @@ component (`Components.DockerProxy`). ``` runner --[registry-mirror]--> zot (:5000) --> DockerHub (first pull of a tag/blob only) - | + | \ + | \--307--> S3 (blob bytes; runner pulls direct) v S3 (manifests + blobs) ``` +With `redirectBlobURL` enabled, blob (layer) pulls are answered with a `307` to a +presigned S3 URL, so the runner downloads the bytes straight from S3 and the zot +NIC carries only manifests and redirects — not the multi-GB layer traffic. + ## Why zot, one process `zot` serves DockerHub images at their **native** `library/...` paths, so Docker's @@ -29,7 +34,12 @@ removes that revalidation, so a single process does the whole job. check, `zot` serves it from S3 and DockerHub is **not** contacted. Otherwise `zot` fetches it from DockerHub once, stores it, and serves it. 2. Runner requests each blob by digest. Blobs are content-addressed, so once in S3 - they are served from S3 and never re-fetched. + they are served from S3 and never re-fetched. With `redirectBlobURL` on, zot + returns a `307` to a presigned S3 URL for the blob (for plain, non-`Range` GETs + — the normal first pull) and the runner fetches the bytes directly from S3; + zot's own NIC is not in the blob data path. If the driver returns no redirect + URL, or the request carries a `Range` header (a resumed download), zot falls + back to streaming the blob through itself. DockerHub is contacted only on the first pull of a tag/blob, plus one manifest re-check per tag after a process restart (the last-check time is in memory). With @@ -70,8 +80,8 @@ Key fields: | Field | Default | Notes | |---|---|---| -| `instance_type` | `c7g.large` | Graviton; the family must end in `g` for the launch-template AMI resolver to pick arm64 (it misses the `gn`/`gd` variants). | -| `zot_version` | `v2.1.21` | Minimum — `manifestCheckInterval` was added here. | +| `instance_type` | `c7g.large` | Graviton (arm64). The AMI resolver keys off a `g` right after the generation number, so network/disk-optimized variants like `c7gn`/`c7gd` resolve correctly too. | +| `zot_version` | `v2.1.21` | Minimum — both `manifestCheckInterval` and `redirectBlobURL` are needed. | | `s3_bucket` / `s3_region` | — | Mirror storage. Reused across instance replacements, so a fresh node boots warm. | | `s3_rootdirectory` | `/zot` | S3 key prefix; isolates zot's layout from anything else in the bucket. | | `manifest_check_interval` | `168h` | Window a cached tag is served without re-checking DockerHub. | @@ -79,6 +89,7 @@ Key fields: | `dns_zone` / `dns_record` | — | Private zone + record the instance self-registers and runners point their registry-mirror at. | | `listen_port` | `5000` | Registry port; runners mirror to `http://:`. | | `enable_ui` | `False` | Serve zot's web UI at `/` (see below). | +| `redirect_blob_url` | `True` | `307`-redirect blob pulls to presigned S3 URLs so runners fetch layers directly from S3 (needs `zot_version` >= v2.1.21). Keeps the proxy NIC out of the multi-GB layer data path. | `dns_zone`, `dns_record`, `s3_bucket` and `dockerhub_pat_ssm` are external identities and are **not** project-namespaced; the IAM role, profile, launch @@ -140,6 +151,12 @@ Settings that must be present (see `docker_proxy_user_data.sh`): - `storage.storageDriver.name: s3` (S3 accessed via the EC2 instance role); `dedupe: false` (no cache driver needed for a single instance); `gc: true` (safe with a single writer). +- `storage.redirectBlobURL: true` (added in zot v2.1.21) — blob GETs return a `307` + to the S3 signed URL so runners pull layer bytes directly from S3. Without it zot + streams every blob byte `S3 -> zot -> runner`, doubling the instance's NIC load; + a single small instance then saturates its EC2 network allowance under CI + fan-out and pulls time out. zot soft-fails back to proxying if no redirect URL is + available, so it is safe to leave on. Under high concurrency, `reqConcurrent` / `reqPerSec` (per-host upstream caps) and `disableHTTP2` are the knobs if DockerHub throttles. diff --git a/praktika/infrastructure/native/docker_proxy.py b/praktika/infrastructure/native/docker_proxy.py index bde0a1e1..0e565dc9 100644 --- a/praktika/infrastructure/native/docker_proxy.py +++ b/praktika/infrastructure/native/docker_proxy.py @@ -63,9 +63,9 @@ class DockerProxy: """ name: str = "dockerhub-proxy" - # Graviton (arm64). Keep a family ending in "g" so the launch-template AMI - # resolver detects arm64 (it keys off family.endswith("g"), which misses the - # "gn"/"gd" network/disk-optimized variants). + # Graviton (arm64). The launch-template AMI resolver detects arm64 from a "g" + # right after the generation number, so network/disk-optimized variants like + # "c7gn"/"c7gd" resolve correctly too. instance_type: str = "c7g.large" vpc_name: str = "" ami_id: str = "" # AL2023 arm64/x86_64 resolved at deploy time if empty @@ -91,6 +91,12 @@ class DockerProxy: # Serve zot's web UI (repo/tag browser) on the same port at `/` (the registry # API stays at `/v2/`). Lightweight: same binary, no CVE/trivy scanning. enable_ui: bool = False + # Set zot's storage.redirectBlobURL so blob (layer) GETs return a 307 to a + # presigned S3 URL and runners download layer bytes straight from S3, keeping + # the proxy NIC out of the data path (only manifests + redirects flow through + # it). Without it zot streams every blob byte S3 -> zot -> runner, doubling the + # instance's NIC load and saturating it under CI fan-out. Needs zot >= v2.1.21. + redirect_blob_url: bool = True # Route53 self-registration. The hosted zone must already exist; the record # is what runners point their registry-mirror at. Both are external DNS # identities and are NOT project-namespaced. @@ -232,6 +238,7 @@ def _refresh(self): dns_record=self.dns_record, tailscale=ts, enable_ui=self.enable_ui, + redirect_blob_url=self.redirect_blob_url, ) def _tailscale_config(self): diff --git a/praktika/infrastructure/native/docker_proxy_user_data.sh b/praktika/infrastructure/native/docker_proxy_user_data.sh index c5c22d8d..7d7d8406 100644 --- a/praktika/infrastructure/native/docker_proxy_user_data.sh +++ b/praktika/infrastructure/native/docker_proxy_user_data.sh @@ -3,15 +3,19 @@ # # Runs a single zot instance as a DockerHub pull-through cache backed by S3: # runner --[registry-mirror]--> zot (:PORT) --> DockerHub (first pull only) -# | +# | \ +# | \--307--> S3 (blob bytes; runner pulls direct) # v # S3 (manifests + blobs) # # zot serves DockerHub images at their native library/... paths, so Docker's # registry-mirror works with no image-reference rewrites. manifestCheckInterval -# makes cached tags serve from S3 without re-contacting DockerHub. The DockerHub -# PAT is read from SSM at boot; S3 uses the EC2 instance role. No static -# credentials are written to disk except the short-lived sync-auth.json (0600). +# makes cached tags serve from S3 without re-contacting DockerHub. With +# redirectBlobURL enabled, blob (layer) GETs return a 307 to a presigned S3 URL, +# so runners download layer bytes straight from S3 and the proxy NIC stays out of +# the data path (only manifests + redirects flow through it). The DockerHub PAT is +# read from SSM at boot; S3 uses the EC2 instance role. No static credentials are +# written to disk except the short-lived sync-auth.json (0600). set -xeuo pipefail # --- Resolve region + private IP from IMDSv2 --- @@ -44,6 +48,7 @@ cat > /etc/zot/config.json <<'ZOTCONF' "rootDirectory": "/var/lib/zot", "dedupe": false, "gc": true, + "redirectBlobURL": __REDIRECT_BLOB_URL__, "storageDriver": { "name": "s3", "region": "__S3_REGION__", diff --git a/praktika/infrastructure/native/user_data.py b/praktika/infrastructure/native/user_data.py index 8dd25588..2f8bf2f5 100644 --- a/praktika/infrastructure/native/user_data.py +++ b/praktika/infrastructure/native/user_data.py @@ -144,6 +144,7 @@ def docker_proxy_user_data( dns_zone, dns_record, enable_ui=False, + redirect_blob_url=True, tailscale=None, ): """Render the DockerHub proxy bootstrap script. @@ -156,6 +157,12 @@ def docker_proxy_user_data( ``enable_ui`` adds zot's ``search`` + ``ui`` extensions (served at ``/`` on the same port). No CVE/trivy scanning is enabled, so it stays lightweight. + ``redirect_blob_url`` (default True) sets zot's ``storage.redirectBlobURL`` so + blob GETs return a ``307`` to the storage driver's signed URL (direct S3 + download) instead of streaming the bytes through zot. Requires zot >= v2.1.21. + zot falls back to proxying if the driver returns no redirect URL, so it is safe + to leave on. + ``tailscale`` (a dict with ``hostname``, ``tag``, ``oauth_client_id_ssm``, ``oauth_client_secret_ssm``) opts the node into Tailscale: it mints a tagged ephemeral auth key from the SSM OAuth client, joins the tailnet with SSH, and @@ -181,6 +188,7 @@ def docker_proxy_user_data( "__DNS_ZONE__": str(dns_zone), "__DNS_RECORD__": str(dns_record), "__EXTRA_EXTENSIONS__": extra_extensions, + "__REDIRECT_BLOB_URL__": "true" if redirect_blob_url else "false", "__TAILSCALE_SETUP__": tailscale_setup, } for placeholder in replacements: diff --git a/praktika/mangle.py b/praktika/mangle.py index 153cc0b6..113563fe 100644 --- a/praktika/mangle.py +++ b/praktika/mangle.py @@ -43,6 +43,7 @@ def _get_workflows( file=None, _for_validation_check=False, _file_names_out=None, + _load_errors_out=None, default=False, ) -> List[Workflow.Config]: """ @@ -84,7 +85,26 @@ def _get_workflows( assert spec foo = importlib.util.module_from_spec(spec) assert spec.loader - spec.loader.exec_module(foo) + try: + spec.loader.exec_module(foo) + except Exception as e: + # One workflow file that fails to import must not take down the whole + # scan (and with it every other workflow). This notably happens on a + # version-skewed pool: a workflow file references a newer praktika + # feature (e.g. a new Workflow.Engine member) than the praktika the + # orchestrator is running, so importing it raises AttributeError. Skip + # the file with a clear warning so the remaining workflows still load + # and match. Validation re-raises so `praktika validate` (run under the + # checkout's own praktika) still catches genuinely broken files loudly. + if _for_validation_check: + raise + print( + f"WARNING: skipping workflow file [{py_file.name}] — failed to " + f"import: {type(e).__name__}: {e}" + ) + if isinstance(_load_errors_out, list): + _load_errors_out.append((py_file.name, f"{type(e).__name__}: {e}")) + continue try: matched_default_workflow = False for workflow in foo.WORKFLOWS: @@ -113,6 +133,12 @@ def _get_workflows( # f"WARNING: Failed to add WORKFLOWS config from [{module_name}], exception [{e}]" # ) if not res: + if isinstance(_load_errors_out, list) and _load_errors_out: + # Every workflow file failed to import (e.g. a pool on an older baked + # praktika than the workflows use). Return empty instead of raising an + # opaque "no workflow found" so the caller can surface the collected + # load errors (and finalize the bootstrap check) rather than crash. + return res Utils.raise_with_error(f"Failed to find [{name or file or 'any'}] workflow") if not _for_validation_check: diff --git a/praktika/orchestrator/__init__.py b/praktika/orchestrator/__init__.py index 891eab97..f2ff8843 100644 --- a/praktika/orchestrator/__init__.py +++ b/praktika/orchestrator/__init__.py @@ -46,8 +46,12 @@ def _current_orchestrator_filter() -> str: return "default" -def find_workflows_for_event(event, workflow_name=None): - """Find all workflows matching the trigger event. Returns empty list if no match.""" +def find_workflows_for_event(event, workflow_name=None, _load_errors_out=None): + """Find all workflows matching the trigger event. Returns empty list if no match. + + ``_load_errors_out``: optional list; workflow files that fail to import are + skipped (so one broken file can't crash the whole scan) and appended here as + ``(filename, error)`` so the caller can surface them (e.g. a GitHub check).""" # Two distinct sources of a name, kept separate on purpose: # - explicit_name: the trusted CLI selector (`orchestrate --name`). It also # bypasses pool (orchestrator_filter) routing, since an operator targeting @@ -80,7 +84,7 @@ def find_workflows_for_event(event, workflow_name=None): matched = [] orchestrator_filter = _current_orchestrator_filter() - for wf in _get_workflows(): + for wf in _get_workflows(_load_errors_out=_load_errors_out): if name_filter and wf.name != name_filter: continue if wf.engine == Workflow.Engine.GH_ACTIONS: @@ -433,12 +437,59 @@ def _orchestrate_event( check = CheckRun(gh_token, event.get("repo", ""), bootstrap_check_id, "CI") - workflows = find_workflows_for_event(event, workflow_name=workflow_name) + load_errors: list = [] + workflows = find_workflows_for_event( + event, workflow_name=workflow_name, _load_errors_out=load_errors + ) + # A workflow file that fails to import is skipped (so one broken file can't + # crash the whole scan), but the skip must not be silent — a developer whose + # workflow didn't run needs to see why. Surface each failed file as its own + # completed "failure" check on the PR head, independent of whether other + # workflows matched (typical cause: a pool pinned to an older baked praktika + # importing a workflow that uses a newer feature). See mangle._get_workflows. + if load_errors and ci and gh_token: + head_sha = event.get("head_sha", "") + repo = event.get("repo", "") + if head_sha and repo: + from .check_run import CheckRun + + for filename, err in load_errors: + try: + # start() opens the check (in_progress), complete() flips it to + # the terminal failure — CheckRun has no one-shot create helper. + # No Cancel action: there's nothing running to cancel. + CheckRun.start( + gh_token, + repo, + head_sha, + f"Workflow load error: {filename}", + with_cancel_action=False, + ).complete( + "failure", + output={ + "title": "Workflow failed to load", + "summary": ( + f"`ci/workflows/{filename}` could not be imported and " + f"was skipped, so any workflow it defines did not run:\n\n" + f"```\n{err}\n```\n\n" + f"This usually means this orchestrator's praktika is " + f"older than a feature the workflow uses." + ), + }, + ) + except Exception as e: + print(f" [warn] could not post workflow-load-error check: {e}") + if not workflows: + # A genuine no-match is neutral; but if workflow files failed to import + # (load_errors) the bootstrap check must not close green/neutral — the + # per-file "Workflow load error" checks above already show red, and this + # completes (never leaves in_progress) the adopted bootstrap check too. print("No matching workflows, exiting") if check is not None: + conclusion = "failure" if load_errors else "neutral" try: - check.complete("neutral", output=_check_output(None, None)) + check.complete(conclusion, output=_check_output(None, None)) except Exception: print(f"Failed to complete check run: {check}", file=sys.stderr) return 0 diff --git a/pyproject.toml b/pyproject.toml index 50ec2d8d..5839dc93 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -4,7 +4,7 @@ build-backend = "setuptools.build_meta" [project] name = "praktika" -version = "0.1.13" +version = "0.1.14" description = "CI Infrastructure Toolbox" requires-python = ">=3.8" license = "Apache-2.0"