#!/usr/bin/env python3 """宿主专用的 Rust 缓存维护器;仅使用 Python 标准库,不在 CI job 中运行。""" import argparse import datetime import hashlib import json import os from pathlib import Path import re import socket import shutil import subprocess import sys import tempfile import time import urllib.error import urllib.parse import urllib.request from gitea_cache_snapshot import ArtifactIdentity, ArtifactInput, merge_snapshots from gitea_cache_upload_cleanup import cleanup_upload_chunks IMAGE = re.compile(r"sha256:[0-9a-f]{64}\Z") SHA = re.compile(r"[0-9a-f]{40}\Z") LABEL = re.compile(r'(?m)^(\s*-\s*[\"\x27]?)genarrative-ci:docker://(sha256:[0-9a-f]{64})([\"\x27]?\s*)$') JOBS = { "AI game creator shell Rust lane 1/2", "AI game creator shell Rust lane 2/2", "AI game creator shell Rust smoke", "AI game creator shell Rust crates", "Backend tests", "Native shell tests", "Frontend tests", "Repository checks", "AI game creator shell web tests", } RUST_JOB_IDS = { "AI game creator shell Rust lane 1/2": "ai-game-creator-shell-rust-lane-1", "AI game creator shell Rust lane 2/2": "ai-game-creator-shell-rust-lane-2", "AI game creator shell Rust smoke": "ai-game-creator-shell-rust-smoke", "AI game creator shell Rust crates": "ai-game-creator-shell-rust-crates", "Backend tests": "backend-tests", "Native shell tests": "native-shell-tests", } RUST_JOBS = set(RUST_JOB_IDS) ARTIFACT_PREFIX = "rust-cache-v1-" MAX_DOWNLOAD = 4 * 1024 ** 3 + 129 * 1024 ** 2 EXPORT_STEP = "Publish master Rust cache artifact" def now(): return datetime.datetime.now(datetime.timezone.utc).isoformat() def timestamp(value): return datetime.datetime.fromisoformat(value.replace("Z", "+00:00")) def log(message): print(f"[cache-maintenance] {message}", flush=True) def command(*args, cwd=None, data=None, env=None, output=None, timeout=120, combined=False): result = subprocess.run( args, cwd=cwd, input=data, text=True, env=env, timeout=timeout, stdout=output or subprocess.PIPE, stderr=output or (subprocess.STDOUT if combined else subprocess.PIPE), ) if result.returncode: # 参数、stderr 或 HTTP body 可能含凭据,只把详细构建输出写到私有日志。 raise RuntimeError(f"{Path(args[0]).name} failed (exit {result.returncode})") return result.stdout or "" def atomic_json(path, value): temporary = path.with_suffix(".tmp") with temporary.open("w", encoding="utf-8") as out: json.dump(value, out, ensure_ascii=False, indent=2) out.write("\n") out.flush() os.fsync(out.fileno()) temporary.replace(path) if os.name == "posix": directory = os.open(path.parent, os.O_RDONLY | os.O_DIRECTORY) try: os.fsync(directory) finally: os.close(directory) def cache_inputs(tree): """只排除已知不参与构建的说明文档;docs/openapi 与内嵌 skill 必须参与。""" rows = [] for row in (tree.split("\0") if "\0" in tree else tree.splitlines()): if not row: continue path = row.split("\t", 1)[1] if path.startswith("docs/") and not path.startswith("docs/openapi/"): continue if path in {"AGENTS.md", "README.md", "deploy/container/README.md"}: continue rows.append(row) return hashlib.sha256("\n".join(rows).encode()).hexdigest() def configured_image(config): labels = list(LABEL.finditer(config)) if len(labels) != 1: raise RuntimeError("expected exactly one genarrative-ci Image ID label") return labels[0].group(2) def replace_image(config, old, new): if not IMAGE.fullmatch(new) or configured_image(config) != old: raise RuntimeError("runner image changed outside maintenance; refusing overwrite") return LABEL.sub(lambda match: match[1] + "genarrative-ci:docker://" + new + match[3], config) def fully_passed(run, jobs): return ( run.get("status") == "completed" and run.get("conclusion") == "success" and len(jobs) == len(JOBS) and {job["name"] for job in jobs} == JOBS and all(job.get("conclusion") == "success" for job in jobs) ) def verifies_image(job, contents, image, require_cache=True): if not re.search(r"(?m)^\S+ image: " + re.escape(image) + r"\s*$", contents): return False if not require_cache or job["name"] not in RUST_JOBS: return True # 旧分支使用了镜像却未 prepare,不算完成缓存上线验收。 hits = re.search(r"Cache hits\s+(\d+)\s*$", contents, re.M) errors = [re.search(rf"{name}\s+(\d+)\s*$", contents, re.M) for name in ("Cache errors", "Cache read errors", "Cache write errors")] return ("[rust-cache] mode=sccache" in contents and hits is not None and int(hits[1]) > 0 and all(item and int(item[1]) == 0 for item in errors)) class Api: def __init__(self, url, token_file): self.url = url.rstrip("/") if not self.url.startswith("https://"): raise ValueError("api_url must use HTTPS") self.token_file = Path(token_file) def request(self, path, *, method, body=None, raw=False): token = self.token_file.read_text().strip() if not token or "\n" in token: raise ValueError("invalid token file") request = urllib.request.Request( self.url + "/" + path.lstrip("/"), data=None if body is None else json.dumps(body).encode(), method=method, headers={"Authorization": "token " + token, "Content-Type": "application/json"}, ) try: # 不允许带 Authorization 的请求跟随跨站重定向。 opener = urllib.request.build_opener(NoRedirect()) with opener.open(request, timeout=30) as response: content = response.read().decode() except urllib.error.HTTPError as error: error.close() if method == "DELETE" and error.code == 404: return None raise RuntimeError(f"Gitea API HTTP {error.code}") from None return content if raw else (json.loads(content) if content else None) def download(self, path, destination): """REST V4 archive redirects to a signed URL; never forward the API token.""" token = self.token_file.read_text().strip() url = self.url + "/" + path.lstrip("/") opener = urllib.request.build_opener(NoRedirect()) request = urllib.request.Request(url, headers={"Authorization": "token " + token}) try: response = opener.open(request, timeout=60) except urllib.error.HTTPError as error: error.close() if error.code not in (301, 302, 303, 307, 308): raise RuntimeError(f"artifact download HTTP {error.code}") from None target = urllib.parse.urljoin(url, error.headers.get("Location", "")) parsed, origin = urllib.parse.urlsplit(target), urllib.parse.urlsplit(self.url) if (parsed.scheme != "https" or parsed.netloc != origin.netloc or parsed.username or parsed.password or target == url): raise RuntimeError("artifact redirect must stay on configured HTTPS Gitea origin") from None response = opener.open(target, timeout=60) try: with response, destination.open("wb") as out: total = 0 while chunk := response.read(1024 * 1024): total += len(chunk) if total > MAX_DOWNLOAD: raise RuntimeError("artifact exceeds per-job size limit") out.write(chunk) except Exception: destination.unlink(missing_ok=True) raise def pages(self, path, key): separator = "&" if "?" in path else "?" page = 1 while True: response = self.request(f"{path}{separator}limit=50&page={page}", method="GET") items = response[key] yield from items if len(items) < 50: return page += 1 class NoRedirect(urllib.request.HTTPRedirectHandler): def redirect_request(self, req, fp, code, msg, headers, newurl): return None class Maintenance: def __init__(self, config): self.config = config if not Path(config["state_dir"]).is_absolute(): raise ValueError("state_dir must be a dedicated absolute directory") self.root = Path(config["state_dir"]).resolve() if self.root == Path("/"): raise ValueError("state_dir must be a dedicated absolute directory") self.repo = self.root / "source" self.runner = config.get("runner_container", "gitea-runner") self.api = Api(config["api_url"], config["token_file"]) self.repo_api = "repos/" + config["repository"] self.state_path = self.root / "state.json" self.state = json.loads(self.state_path.read_text()) if self.state_path.exists() else { "versions": [], "current": None, "rollback": None, "candidate": None, } def save(self): atomic_json(self.state_path, self.state) def docker(self, *args, inner=False, **kwargs): prefix = ("docker", "exec", self.runner, "docker") if inner else ("docker",) return command(*prefix, *args, **kwargs) def image_info(self, image, inner=False): return json.loads(self.docker("image", "inspect", image, inner=inner))[0] def read_config(self): return self.docker("exec", self.runner, "cat", "/data/config.yaml") def write_config(self, content): self.docker("exec", "--user", "0", "-i", self.runner, "sh", "-c", "cat > /data/config.yaml", data=content) def version(self, image): return next(row for row in self.state["versions"] if row["image"] == image) def adopt_current(self): actual = configured_image(self.read_config()) if self.state["current"]: if actual != self.state["current"]: raise RuntimeError("runner configuration differs from recorded current image") return info = self.image_info(actual) labels = info["Config"].get("Labels") or {} source = labels.get("world.genarrative.ci.rust-cache-source") base = labels.get("world.genarrative.ci.rust-cache-base", actual) if source is not None and not SHA.fullmatch(source): raise RuntimeError("invalid image source SHA") self.state["versions"].append({ "image": actual, "source": source, "base": base, "owned": False, "activated": "1970-01-01T00:00:00+00:00", "verified_run": None, }) self.state["current"] = actual self.save() def fetch_source(self): if self.repo.resolve() != self.repo: raise RuntimeError("dedicated clone must not be a symlink") marker = self.repo / ".git" / "genarrative-cache-maintenance" if not (self.repo / ".git").exists(): if self.repo.exists(): raise RuntimeError("source directory already exists without .git") command("git", "clone", "--no-checkout", self.config["clone_url"], str(self.repo), timeout=900) marker.write_text("owned source checkout\n") if not marker.is_file(): raise RuntimeError("refusing to clean a checkout not created by this maintainer") if command("git", "remote", "get-url", "origin", cwd=self.repo).strip() != self.config["clone_url"]: raise RuntimeError("dedicated clone origin differs from configured repository") command("git", "fetch", "--no-tags", "origin", "refs/heads/master", cwd=self.repo, timeout=900) sha = command("git", "rev-parse", "FETCH_HEAD^{commit}", cwd=self.repo).strip() if not SHA.fullmatch(sha): raise RuntimeError("invalid master SHA") # 这是维护器专属 clone,不修改开发者工作区;编译期间不再 checkout 新提交。 command("git", "checkout", "--detach", "--force", sha, cwd=self.repo) command("git", "clean", "-ffdx", cwd=self.repo) tree = command("git", "ls-tree", "-rz", sha, cwd=self.repo) return sha, cache_inputs(tree) def build_command(self, log_file, script, *args, env=None): with log_file.open("a") as out: command("bash", str(self.repo / "scripts" / script), *args, cwd=self.repo, env=env, output=out, timeout=7200) def master_run(self, run): return (run.get("path") == "project-ci.yml@refs/heads/master" and run.get("event") == "push" and run.get("head_branch") == "master" and SHA.fullmatch(run.get("head_sha", "")) is not None) def source_run(self): """Only complete exports from the latest eligible master run; never mix runs.""" current = self.version(self.state["current"]) current_inputs = (cache_inputs(command("git", "ls-tree", "-rz", current["source"], cwd=self.repo)) if current.get("source") else None) for run in self.api.pages(self.repo_api + "/actions/runs?branch=master&event=push", "workflow_runs"): if (not self.master_run(run) or run.get("status") != "completed" or run.get("conclusion") not in {"success", "failure"} or run["id"] <= current.get("run_id", 0)): continue sha = run["head_sha"] if sha == current.get("source"): continue ancestry = subprocess.run(["git", "merge-base", "--is-ancestor", sha, "FETCH_HEAD"], cwd=self.repo, stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL) if ancestry.returncode: continue if current.get("source"): ancestry = subprocess.run(["git", "merge-base", "--is-ancestor", current["source"], sha], cwd=self.repo, stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL) if ancestry.returncode: continue inputs = cache_inputs(command("git", "ls-tree", "-rz", sha, cwd=self.repo)) if inputs == current_inputs: continue jobs = list(self.api.pages(self.repo_api + f'/actions/runs/{run["id"]}/jobs', "jobs")) rust_jobs = [job for job in jobs if job["name"] in RUST_JOB_IDS] if len(rust_jobs) != len(RUST_JOB_IDS) or {job["name"] for job in rust_jobs} != RUST_JOBS: continue if any(job.get("status") != "completed" or job.get("conclusion") not in {"success", "failure"} or job.get("head_sha") != sha or not any(step.get("name") == EXPORT_STEP and step.get("conclusion") == "success" for step in job.get("steps", [])) for job in rust_jobs): continue artifacts = list(self.api.pages(self.repo_api + f'/actions/runs/{run["id"]}/artifacts', "artifacts")) selected = [] images = set() for job in rust_jobs: name = f'{ARTIFACT_PREFIX}{RUST_JOB_IDS[job["name"]]}-attempt-{job["run_attempt"]}' matches = [item for item in artifacts if item["name"] == name and not item.get("expired") and item.get("workflow_run", {}).get("id") == run["id"]] if len(matches) != 1: break content = self.api.request(self.repo_api + f'/actions/jobs/{job["id"]}/logs', method="GET", raw=True) if not re.search(re.escape(f"[rust-cache] artifact={name}") + r" objects=\d+ bytes=\d+ complete=true", content): break used = set(re.findall(r"(?m)^\S+ image: (sha256:[0-9a-f]{64})\s*$", content)) if len(used) != 1: break images.update(used) selected.append({"id": matches[0]["id"], "name": name, "job": RUST_JOB_IDS[job["name"]], "attempt": job["run_attempt"]}) if len(selected) != len(RUST_JOB_IDS) or len(images) != 1: continue source_image = images.pop() if source_image not in {row["image"] for row in self.state["versions"]}: continue return {"run_id": run["id"], "source": sha, "inputs": inputs, "source_image": source_image, "exports": selected} return None def build(self, source): """Assemble existing CI objects; no cargo warm-up or test execution.""" sha = source["source"] command("git", "checkout", "--detach", "--force", sha, cwd=self.repo) command("git", "clean", "-ffdx", cwd=self.repo) artifact = self.root / "artifacts" / sha artifact.mkdir(parents=True, exist_ok=True) build_log = artifact / "build.log" tag = "genarrative/gitea-project-ci:rust-cache-auto-" + sha base_tag = "genarrative/gitea-project-ci:base-auto-" + sha attempt = {**source, "tag": tag, "base_tag": base_tag, "artifact": str(artifact)} self.state.setdefault("attempts", []).append(attempt) self.save() env = {**os.environ, "GENARRATIVE_GITEA_RUNNER_CONTAINER": self.runner} env.pop("CI", None) labels = self.image_info(source["source_image"])["Config"].get("Labels") or {} inherited_source = labels.get("world.genarrative.ci.rust-cache-source") base = labels.get("world.genarrative.ci.rust-cache-base") if not IMAGE.fullmatch(base or "") or not SHA.fullmatch(inherited_source or ""): raise RuntimeError("source image must contain a trusted cache snapshot") revision = command("bash", "scripts/gitea-ci-job-image.sh", "revision", cwd=self.repo).strip() base_labels = self.image_info(base)["Config"].get("Labels") or {} if base_labels.get("com.genarrative.ci.definition-sha256") != revision: env["GENARRATIVE_GITEA_CI_IMAGE_TAG"] = base_tag self.build_command(build_log, "gitea-ci-job-image.sh", "build", env=env) base = self.image_info(base_tag)["Id"] self.state.setdefault("bases", {})[base] = base_tag self.save() # 与旧缓存镜像分离;绝不把 Docker 可写层、源码或 target commit 成镜像。 self.docker("run", "--rm", "--network", "none", "--read-only", "--cap-drop=ALL", "--entrypoint", "bash", base, "-c", "test ! -e /opt/genarrative-ci/rust-cache") with tempfile.TemporaryDirectory(prefix="assemble-", dir=artifact) as temporary: work = Path(temporary) inherited = work / "inherited" inherited.mkdir() container = self.docker("create", source["source_image"]).strip() try: self.docker("cp", container + ":/opt/genarrative-ci/rust-cache/.", str(inherited), timeout=600) finally: self.docker("rm", "--volumes", container) inputs = [] for export in source["exports"]: archive = work / (str(export["id"]) + ".zip") self.api.download(self.repo_api + f'/actions/artifacts/{export["id"]}/zip', archive) inputs.append(ArtifactInput(archive, ArtifactIdentity( self.config["repository"], source["run_id"], export["attempt"], export["job"], sha))) snapshot = work / "snapshot" merged = merge_snapshots(inputs, snapshot, base_objects=inherited / "objects", expected_inherited_source_sha=inherited_source) if merged.sccache_version != "sccache 0.18.0": raise RuntimeError("unsupported sccache version") if merged.base_image is not None and merged.base_image != labels["world.genarrative.ci.rust-cache-base"]: raise RuntimeError("artifact base differs from its actual source image") rustc = self.docker("run", "--rm", "--network", "none", "--read-only", base, "rustc", "-vV") if rustc.strip() != merged.rustc.strip(): raise RuntimeError("artifact toolchain differs from target base image") if merged.workspace != "/workspace/" + self.config["repository"]: raise RuntimeError("artifact workspace differs from CI checkout") shutil.copyfile(inherited / "sccache", snapshot / "sccache") (snapshot / "sccache").chmod(0o755) (snapshot / "base-image.txt").write_text(base + "\n") (work / "Dockerfile").write_text( f"FROM {base}\nCOPY snapshot/ /opt/genarrative-ci/rust-cache/\n" f'LABEL world.genarrative.ci.rust-cache-source="{sha}"\n' f'LABEL world.genarrative.ci.rust-cache-base="{base}"\n') (work / ".dockerignore").write_text("**\n!Dockerfile\n!snapshot/\n!snapshot/**\n") with build_log.open("a") as out: self.docker("build", "--pull=false", "--tag", tag, str(work), output=out, timeout=1800) self.build_command(build_log, "gitea-ci-job-image.sh", "verify", tag, env=env) image = self.image_info(tag)["Id"] self.state["versions"].append({**attempt, "image": image, "base": base, "owned": True, "verified_run": None}) self.state["candidate"] = image self.state["attempts"] = [row for row in self.state["attempts"] if row["source"] != sha] self.save() log(f"candidate assembled run={source['run_id']} source={sha} image={image}") def stage_candidate(self): candidate = self.version(self.state["candidate"]) command("git", "checkout", "--detach", "--force", candidate["source"], cwd=self.repo) command("git", "clean", "-ffdx", cwd=self.repo) artifact = Path(candidate["artifact"]) archive = artifact / "image.tar.zst" sidecar = archive.with_suffix(".zst.sha256") if (candidate.get("staged") and archive.is_file() and sidecar.is_file() and candidate["image"] in self.docker("image", "ls", "--all", "--no-trunc", "--quiet", inner=True).split()): return env = {**os.environ, "GENARRATIVE_GITEA_RUNNER_CONTAINER": self.runner} build_log = artifact / "build.log" if not sidecar.exists(): # 只删除登记目录中的未完成导出文件,不覆盖已验证归档。 archive.unlink(missing_ok=True) self.build_command(build_log, "gitea-ci-job-image.sh", "export", str(archive), candidate["image"], env=env) command("sha256sum", "--check", sidecar.name, cwd=artifact, timeout=600) self.build_command(build_log, "gitea-ci-job-image.sh", "load-runner", candidate["image"], env=env) if self.image_info(candidate["image"], inner=True)["Id"] != candidate["image"]: raise RuntimeError("inner runner image mismatch") candidate["staged"] = True self.save() def idle(self): gate = self.gate("status") if type(gate.get("active_tasks")) is not int or gate["active_tasks"] < 0: raise RuntimeError("gate must support durable task tracking before automatic switching") # 已领取但尚未建容器、正在收尾上报的任务都由入口跟踪,不查管理员 API。 return (gate.get("uncertain") is False and gate.get("active_tasks") == 0 and not self.docker("ps", "-q", "--filter", "status=running", "--filter", "status=created", "--filter", "status=restarting", "--filter", "status=paused", inner=True).strip()) def verify_current(self): current = self.version(self.state["current"]) if current.get("verified_run"): return True for run in self.api.pages(self.repo_api + "/actions/runs?status=success&branch=master&event=push", "workflow_runs"): if run.get("path") != "project-ci.yml@refs/heads/master" or run.get("event") != "push": continue sha = run.get("head_sha", "") if not SHA.fullmatch(sha): continue ancestry = subprocess.run(["git", "merge-base", "--is-ancestor", sha, "FETCH_HEAD"], cwd=self.repo, stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL) if ancestry.returncode != 0: continue if timestamp(run["started_at"]) < timestamp(current["activated"]): continue jobs = list(self.api.pages(self.repo_api + f'/actions/runs/{run["id"]}/jobs', "jobs")) if not fully_passed(run, jobs): continue if all(verifies_image(job, self.api.request( self.repo_api + f'/actions/jobs/{job["id"]}/logs', method="GET", raw=True), current["image"], require_cache=current.get("source") is not None, ) for job in jobs): current["verified_run"] = run["id"] current["verified_sha"] = sha self.save() log(f'current image verified by run={run["id"]}') return True log("waiting for a complete successful CI run using current image and cache") return False def protected_images(self): protected = {self.state[key] for key in ("current", "rollback", "candidate") if self.state.get(key)} for image in list(protected): protected.add(self.version(image)["base"]) # 保护包括已停止容器在内的全部引用;不强制删除被 Docker 引用的镜像。 for inner in (False, True): ids = self.docker("ps", "-aq", inner=inner).split() for container in ids: protected.add(self.docker("inspect", "--format", "{{.Image}}", container, inner=inner).strip()) protected.add(configured_image(self.read_config())) return protected def remove_owned_image(self, image, tag, protected): if image in protected or not IMAGE.fullmatch(image): return False for inner in (True, False): available = self.docker("image", "ls", "--all", "--no-trunc", "--quiet", inner=inner).split() if image not in available: continue info = self.image_info(image, inner=inner) if not set(info.get("RepoTags") or []).issubset({tag}): log(f"retain image with additional tags: {image}") return False # 不使用 --force;并发创建的容器也由 Docker 的引用检查保护。 self.docker("image", "rm", image, inner=inner) if image in self.docker("image", "ls", "--all", "--no-trunc", "--quiet", inner=inner).split(): raise RuntimeError("managed image is still present after removal") return True def cleanup(self): if not self.version(self.state["current"]).get("verified_run"): return protected = self.protected_images() for record in list(self.state["versions"]): if not record.get("owned") or record["image"] in protected: continue expected = self.root / "artifacts" / record["source"] if (not SHA.fullmatch(record["source"]) or Path(record["artifact"]).resolve() != expected or expected.is_symlink() or expected.parent.is_symlink() or record["tag"] != "genarrative/gitea-project-ci:rust-cache-auto-" + record["source"]): raise RuntimeError("archive path outside managed source directory") if not self.remove_owned_image(record["image"], record["tag"], protected): continue # 不递归删除目录,只删除本维护器创建的已知文件。 for name in ("image.tar.zst", "image.tar.zst.sha256", "build.log"): (expected / name).unlink(missing_ok=True) if expected.exists() and not any(expected.iterdir()): expected.rmdir() self.state["versions"].remove(record) self.save() log(f'removed old managed snapshot {record["image"]}') for image, tag in list(self.state.get("bases", {}).items()): used_bases = {row["base"] for row in self.state["versions"]} if image not in used_bases and self.remove_owned_image(image, tag, protected): del self.state["bases"][image] self.save() def cleanup_exports(self): """Only our named artifacts; keep logs/runs and every unrelated artifact.""" protected_runs = {row["run_id"] for row in self.state.get("attempts", []) if row.get("run_id")} for row in self.state["versions"]: if row.get("run_id") and not row.get("staged"): protected_runs.add(row["run_id"]) if not row.get("staged"): continue for item in row.get("exports", []): if item.get("deleted"): continue self.api.request(self.repo_api + f'/actions/artifacts/{item["id"]}', method="DELETE") item["deleted"] = True self.save() cutoff = datetime.datetime.now(datetime.timezone.utc) - datetime.timedelta(days=7) # 先收集再删除,避免按页删除让下一页位置前移、漏掉旧产物。 artifacts = list(self.api.pages(self.repo_api + "/actions/artifacts", "artifacts")) pattern = re.compile(re.escape(ARTIFACT_PREFIX) + "(" + "|".join(RUST_JOB_IDS.values()) + r")-attempt-\d+\Z") runs = {} for item in artifacts: run_id = (item.get("workflow_run") or {}).get("id") if (not pattern.fullmatch(item["name"]) or not run_id or run_id in protected_runs or timestamp(item["created_at"]) >= cutoff): continue # Artifact.workflow_run 在 Gitea 1.26.4 中只有 id/repository_id/head_sha。 if run_id not in runs: runs[run_id] = self.api.request(self.repo_api + f"/actions/runs/{run_id}", method="GET") run = runs[run_id] if not self.master_run(run) or run.get("status") != "completed": continue self.api.request(self.repo_api + f'/actions/artifacts/{item["id"]}', method="DELETE") def cleanup_pending_uploads(self): # Gitea 1.26.4 的过期/DELETE API 不会清理未 finalized 的 V4 分块。 # 只处理本上传器命名的块,且 run 和文件本身均已过保留期限。 cutoff = datetime.datetime.now(datetime.timezone.utc) - datetime.timedelta(days=7) eligible = set() protected = {row["run_id"] for row in self.state.get("attempts", []) if row.get("run_id")} for run in self.api.pages(self.repo_api + "/actions/runs?branch=master&event=push", "workflow_runs"): if (self.master_run(run) and run.get("status") == "completed" and run["id"] not in protected and run.get("completed_at") and timestamp(run["completed_at"]) < cutoff): eligible.add(run["id"]) removed = cleanup_upload_chunks(Path(self.config["artifact_storage_dir"]), eligible, cutoff.timestamp()) if removed: log(f"removed {removed} expired owned upload fragments") def tick(self, retry=False): if not self.recover_switch(): return self.adopt_current() self.fetch_source() self.recover_builds() self.cleanup_exports() self.cleanup_pending_uploads() if not self.verify_current(): return self.cleanup() if self.state.get("candidate"): self.stage_candidate() self.activate() return source = self.source_run() if source is None: log("waiting for a complete set of master CI cache exports") return if not retry and self.state.get("failed_run") == source["run_id"]: log(f'previous assembly failed at run={source["run_id"]}; waiting for new run or --retry') return # 下载、合并和镜像装载也消耗宿主 IO;繁忙时留给 CI,下轮再收集。 if not self.idle(): log("CI active; defer refresh") return try: self.build(source) except Exception: self.state["failed_run"] = source["run_id"] self.save() raise self.state.pop("failed_run", None) self.save() self.stage_candidate() self.cleanup_exports() self.activate() def recover_builds(self): # 宕机可能发生在 docker build 完成之后、登记 Image ID 之前。 # 只接管预先登记的确定性 tag;恢复的候选仍须经过 stage 的 verify/load。 for attempt in list(self.state.get("attempts", [])): artifact = self.root / "artifacts" / attempt["source"] if (not SHA.fullmatch(attempt["source"]) or Path(attempt["artifact"]).resolve() != artifact or artifact.is_symlink() or artifact.parent.is_symlink()): raise RuntimeError("interrupted assembly directory is outside managed artifacts") # 持有维护锁,只有此前中断的组装可能遗留这些私有工作目录。 for directory in artifact.glob("assemble-*"): if directory.is_symlink() or not directory.is_dir(): raise RuntimeError("unexpected interrupted assembly entry") shutil.rmtree(directory) for kind in ("base_tag", "tag"): ids = set(self.docker("image", "ls", "--no-trunc", "--quiet", attempt[kind]).split()) if not ids: continue if len(ids) != 1: raise RuntimeError("ambiguous managed build tag") image = ids.pop() if kind == "base_tag": self.state.setdefault("bases", {})[image] = attempt[kind] elif not any(row["image"] == image for row in self.state["versions"]): labels = self.image_info(image)["Config"].get("Labels") or {} if labels.get("world.genarrative.ci.rust-cache-source") != attempt["source"]: raise RuntimeError("interrupted build tag has unexpected source") self.state["versions"].append({ **attempt, "image": image, "owned": True, "verified_run": None, "base": labels["world.genarrative.ci.rust-cache-base"], }) if self.state.get("candidate") not in (None, image): raise RuntimeError("interrupted build conflicts with an existing candidate") self.state["candidate"] = image self.state["attempts"].remove(attempt) self.save() def recover_switch(self): if self.state.get("switch"): # 不能仅凭配置文件已经替换,就认定内存中的 runner 已加载新映射。 log("recovering interrupted image switch") return self.activate() if self.state.get("pause_owned"): self.resume() return True def resume(self): if self.state.get("pause_owned"): if self.gate("resume").get("paused") is not False: raise RuntimeError("runner resume could not be confirmed") self.state.pop("pause_owned", None) self.save() def gate(self, action): with socket.socket(socket.AF_UNIX, socket.SOCK_STREAM) as client: client.settimeout(10) client.connect(self.config["gate_socket"]) client.sendall((json.dumps({"action": action}) + "\n").encode()) with client.makefile("rb") as stream: response = json.loads(stream.readline(4096)) if "error" in response: raise RuntimeError("runner gate refused command") return response def check_gate_route(self): checkout = re.search(r"(?m)^\s+GENARRATIVE_GITEA_REPOSITORY_URL:\s*(.+?)\s*$", self.read_config()) if not checkout or checkout[1].strip("\"'") != self.config["repository_url"]: raise RuntimeError("runner.envs must pin GENARRATIVE_GITEA_REPOSITORY_URL independently of the RPC gate") # .runner 含认证材料,只在进程内取 address,不回显原文。 registration = json.loads(self.docker("exec", self.runner, "cat", "/data/.runner")) if registration.get("address", "").rstrip("/") != self.config["gate_url"].rstrip("/"): raise RuntimeError("runner is not configured to use the FetchTask gate") info = json.loads(self.docker("inspect", self.runner))[0] started = timestamp(info["State"]["StartedAt"]).timestamp() addresses = {network["IPAddress"] for network in info["NetworkSettings"]["Networks"].values()} gate = self.gate("status") # .runner 会因 label 更新而回写;mtime 不能证明内存中的地址。 # 要求入口实际见到本次容器启动后、来自它的 FetchTask。 if gate.get("last_fetch_peer") not in addresses or (gate.get("last_fetch_at") or 0) < started: raise RuntimeError("gate has not observed FetchTask from this runner startup") if type(gate.get("active_tasks")) is not int or gate["active_tasks"] < 0: raise RuntimeError("gate must support durable task tracking before automatic switching") def activate(self): self.check_gate_route() candidate = self.version(self.state["candidate"]) config = self.read_config() gate = self.gate("status") if gate.get("paused") is not False and not self.state.get("pause_owned"): log("runner gate paused by operator; defer switch") return False if not self.state.get("switch") and not self.idle(): log("CI active; candidate stays staged") return False # 先持久化恢复意图;控制请求超时也可能已生效,ExecStopPost/下次 tick 会恢复。 self.state["pause_owned"] = True self.save() try: if self.gate("pause").get("paused") is not True: raise RuntimeError("runner pause could not be confirmed") # 不用 FetchTask 客户端超时猜测服务端事务是否已经结束。 # 入口必须完整读完已转发的响应;不确定时拒绝自动重启。 for _ in range(30): gate = self.gate("status") if gate.get("uncertain"): raise RuntimeError("in-flight FetchTask completion is uncertain; manual gate inspection required") if gate.get("paused") is not True: raise RuntimeError("runner gate unexpectedly resumed") if gate.get("inflight") == 0: break time.sleep(1) else: log("FetchTask still in flight; defer switch") return False if not self.idle(): log("in-flight task appeared; defer switch without stopping runner") return False latest_config = self.read_config() if latest_config != config: raise RuntimeError("runner config changed while pausing") pending = self.state.get("switch") if pending is None: old = self.state["current"] updated = replace_image(config, old, candidate["image"]) backup = self.root / "backups" / (candidate["source"] + ".yaml") backup.parent.mkdir(parents=True, exist_ok=True) backup.write_text(config) pending = {"old": old, "new": candidate["image"], "backup": str(backup)} self.state["switch"] = pending self.save() else: actual = configured_image(config) if actual not in {pending["old"], pending["new"]}: raise RuntimeError("interrupted switch conflicts with operator configuration") updated = replace_image(config, actual, pending["new"]) self.write_config(updated) if self.read_config() != updated or not self.idle(): # 配置尚未重启加载,恢复旧文件即可;不干扰意外出现的 job。 self.write_config(Path(pending["backup"]).read_text()) self.state.pop("switch", None) self.save() log("idle check changed before restart; restored configuration") return False self.docker("restart", "--timeout", "660", self.runner, timeout=720) started = self.docker("inspect", "--format", "{{.State.StartedAt}}", self.runner).strip() ready = False for _ in range(30): try: info = self.docker("inspect", "--format", "{{.State.Status}}", self.runner).strip() recent = self.docker("logs", "--since", started, self.runner, combined=True) ready = (info == "running" and "declare successfully" in recent and self.image_info(candidate["image"], inner=True)["Id"] == candidate["image"]) except RuntimeError: ready = False if ready: break time.sleep(2) if not ready: raise RuntimeError("runner registration not confirmed; pending switch retained for recovery") candidate["activated"] = now() self.state["rollback"] = pending["old"] self.state["current"] = pending["new"] self.state["candidate"] = None self.state.pop("switch", None) self.save() log(f'activated {candidate["image"]}; awaiting real master CI validation') return True finally: # 持久化 switch 保证崩溃后不会错误宣称新镜像已验证。 self.resume() def main(): parser = argparse.ArgumentParser(description=__doc__) parser.add_argument("--config", required=True) parser.add_argument("--apply", action="store_true", help="执行维护;默认只检查连接和配置") parser.add_argument("--retry", action="store_true", help="重试同一 master run 上失败的缓存组装") parser.add_argument("--resume", action="store_true", help="仅恢复本维护器暂停的领取;用于 ExecStopPost") args = parser.parse_args() if os.environ.get("CI") == "true": parser.error("run on the trusted host, outside CI jobs") os.umask(0o077) config = json.loads(Path(args.config).read_text()) maintenance = Maintenance(config) if not args.apply and not args.resume: maintenance.api.request(maintenance.repo_api + "/actions/artifacts?limit=1", method="GET") image = configured_image(maintenance.read_config()) head = command("git", "ls-remote", config["clone_url"], "refs/heads/master").split()[0] maintenance.check_gate_route() cleanup_upload_chunks(Path(config["artifact_storage_dir"]), set(), time.time()) gate = maintenance.gate("status") if gate.get("uncertain"): raise RuntimeError("runner gate needs manual inspection") log(f'check master={head} image={image} active_tasks={gate.get("active_tasks")}') return import fcntl # Linux 宿主;纯逻辑测试仍可在 Windows 上导入。 maintenance.root.mkdir(parents=True, exist_ok=True) with (maintenance.root / "maintenance.lock").open("a") as lock: try: fcntl.flock(lock, fcntl.LOCK_EX | fcntl.LOCK_NB) except BlockingIOError: log("another maintenance process holds the lock") return # 加锁后重新读取状态,避免另一个进程完成后覆盖旧快照。 maintenance = Maintenance(config) if args.resume: maintenance.resume() else: maintenance.tick(retry=args.retry) if __name__ == "__main__": try: main() except Exception as error: log(f"maintenance failed: {type(error).__name__}: {error}") sys.exit(1)