Compare commits

...
Author SHA1 Message Date
Satyam Srivastava 0ef47caf86 [bugfix]: select architecture-specific Blackwell kernel targets
- Detect GB200 and GB300 as 10.0a and 10.3a
- Build future GB200 CI images with 10.0a
- Add regression tests for detection, cache keys, and build arguments
2026-10-05 15:42:48 -07:00
Satyam Srivastava e71d01648c [feat]: add selectable Modal and GB200 Kubernetes GPU CI
- Preserve existing CI and add configurable GPU backend routing
- Limit execution to 2 active PRs, 4 GPUs per PR, and 8 GPUs total
- Add persistent admission, cancellation recovery, and backend statuses
- Document deployment and add infrastructure regression tests
2026-10-01 14:44:14 -07:00
25 changed files with 2997 additions and 17 deletions
+6
View File
@@ -22,6 +22,12 @@ else
export PERF_UPLOAD_POLICY=never
fi
# Alternate GPU backends compare against references without publishing records.
# Their worker has read-only Hub credentials; publication is an operator task.
if [ "${FASTVIDEO_CI_LOCAL_ONLY:-0}" = 1 ]; then
export PERF_UPLOAD_POLICY=never
fi
nvidia-smi \
--query-gpu=index,timestamp,clocks.sm,clocks.max.sm,power.draw,power.limit,temperature.gpu \
--format=csv -l 10 > "$PERF_REPORTS_DIR/gpu_telemetry.csv" 2>/dev/null &
@@ -11,6 +11,7 @@ jobs:
if: >-
github.event.context == 'direct-test-completed'
&& github.event.state == 'success'
&& (vars.CI_GPU_BACKEND == '' || vars.CI_GPU_BACKEND == 'slurm')
runs-on: ubuntu-latest
steps:
- name: Check and update aggregate status
@@ -0,0 +1,62 @@
name: Promote Selected GPU Backend Status
on:
status:
permissions:
statuses: write
concurrency:
group: gpu-ci-status-${{ github.event.sha }}-${{ vars.CI_GPU_BACKEND }}
cancel-in-progress: false
jobs:
promote:
if: >-
(vars.CI_GPU_BACKEND == 'modal' || vars.CI_GPU_BACKEND == 'vllm')
&& (github.event.context == format('gpu-ci/{0}/fastcheck-passed', vars.CI_GPU_BACKEND)
|| github.event.context == format('gpu-ci/{0}/full-suite-passed', vars.CI_GPU_BACKEND))
runs-on: ubuntu-latest
env:
SELECTED_BACKEND: ${{ vars.CI_GPU_BACKEND }}
steps:
- name: Mirror the selected backend's latest suite results
uses: actions/github-script@60a0d83039c74a4aee543508d2ffcb1c3799cdea # v7.0.1
with:
script: |
const backend = process.env.SELECTED_BACKEND;
if (!['modal', 'vllm'].includes(backend)) {
throw new Error('Unsupported selected GPU backend');
}
const sha = context.payload.sha;
// Read current state after entering the serialized workflow. A
// delayed event must not overwrite a newer failure with success.
const statuses = await github.paginate(github.rest.repos.listCommitStatusesForRef, {
owner: context.repo.owner,
repo: context.repo.repo,
ref: sha,
per_page: 100,
});
for (const suffix of ['fastcheck-passed', 'full-suite-passed']) {
const sourceContext = `gpu-ci/${backend}/${suffix}`;
const matches = statuses.filter(status => status.context === sourceContext);
matches.sort((a, b) =>
Date.parse(b.updated_at) - Date.parse(a.updated_at) || b.id - a.id
);
const latest = matches[0];
const state = latest ? latest.state : 'pending';
if (!['pending', 'success', 'failure', 'error'].includes(state)) {
throw new Error(`Unsupported status state for ${sourceContext}`);
}
await github.rest.repos.createCommitStatus({
owner: context.repo.owner,
repo: context.repo.repo,
sha,
context: suffix,
state,
description: latest
? `${backend} ${suffix}: ${state}`
: `Waiting for ${backend} ${suffix}`,
...(latest && latest.target_url ? {target_url: latest.target_url} : {}),
});
}
+3 -2
View File
@@ -203,7 +203,8 @@ jobs:
docker buildx imagetools create "${TAG_ARGS[@]}" "${IMAGE_REFS[@]}"
docker buildx imagetools inspect "${TAGS[0]}"
# The CI runner is ARM64 like DGX Spark, but targets sm_100 rather than sm_121.
# The CI runner is ARM64 like DGX Spark, but targets sm_100a rather than sm_121.
# The architecture-specific target includes the GB200 VSA CUDA extensions.
# Publish a single-architecture variant so the self-hosted CI runner can reuse
# the exact prebuilt kernel instead of compiling it in every job.
build-ci-runner-image:
@@ -219,7 +220,7 @@ jobs:
PYTHON_VERSION=3.12
CUDA_VERSION=13.0.0
UV_TORCH_BACKEND=cu130
TORCH_CUDA_ARCH_LIST=10.0
TORCH_CUDA_ARCH_LIST=10.0a
CMAKE_BUILD_PARALLEL_LEVEL=1
FLASH_ATTN_WHEEL_TAG=cu130torch2.12
FLASH_ATTN_WHEEL_RELEASE_ARM64=https://github.com/mjun0812/flash-attention-prebuild-wheels/releases/download/v0.9.22
+10 -2
View File
@@ -4,6 +4,11 @@ This is the canonical reference for FastVideo's CI/CD system. Contributor-facing
PR steps live in [Pull Requests](pull_requests.md), and test-authoring guidance
lives in [Testing](testing.md).
The existing Slurm route below remains the default. Operators can also install
the [selectable GPU dispatcher](gpu_ci_backends.md) to run the same lane scripts
on Modal or Kubernetes in the `vllm` namespace. That opt-in path enforces two
active PRs, four GPUs per PR, and eight total, with separate backend statuses.
## Overview
FastVideo splits validation across GitHub Actions, Buildkite, Slinky Slurm,
@@ -435,8 +440,11 @@ before later jobs consume the updated image pin.
The same workflow publishes a single-architecture ARM64, CUDA 13, SM100 image
for the self-hosted CI runner under the
`py3.12-cuda13.0.0-sm100-{latest,sha-*}` tags. It carries the matching prebuilt
kernel wheel so runner jobs can validate and install the exact source and ABI
match instead of recompiling it in every lane.
kernel wheel compiled with `TORCH_CUDA_ARCH_LIST=10.0a` to include the GB200
VSA CUDA extensions. Runtime kernel detection uses the same target, so runner
jobs can validate and install the exact source and ABI match instead of
recompiling it in every lane. Older artifacts built for `10.0` have a different
cache key and trigger a local rebuild when the worker detects `10.0a`.
The optional Dreamverse matrix builds backend and UI images for CUDA 12.6 and
CUDA 13 on `amd64`. Dreamverse remains `amd64`-only because its FA4 dependency
+252
View File
@@ -0,0 +1,252 @@
# Selectable GPU CI Backends
GPU CI can use the existing Slurm dispatcher, Modal, or GB200 Kubernetes Jobs
in the `vllm` namespace. GitHub and Buildkite remain the trigger and reporting
systems. `vllm` names the cluster namespace here; tests run FastVideo's existing
lane scripts rather than a vLLM inference server.
The default remains Slurm. The existing `.buildkite/pipeline.yml`, private
Slurm uploader, and dormant Modal launchers are preserved. Deploying the new
trusted uploader is an operator step; merging these files does not change the
live Buildkite pipelines or create cluster resources. See
[CI/CD Architecture](ci_architecture.md) for the existing installation.
## Execution and Limits
```text
GitHub PR / slash command / schedule
-> trusted Buildkite bootstrap
-> slurm: existing validated graph and private dispatcher
-> modal or vllm: one trusted suite coordinator
-> shared PR admission and GPU reservations
-> isolated GPU worker per existing lane
-> lane results, logs, and backend-specific suite status
```
All new Modal and `vllm` coordinators share one admission database. It enforces:
- At most two distinct PRs with admitted GPU work.
- At most four reserved GPUs per PR across its concurrent builds and lanes.
- At most eight reserved GPUs across both backends combined.
- Additional PRs wait without creating GPU workers.
A PR is identified by repository and PR number, so Fastcheck, merge builds,
manual reruns, and separate attempts share its allowance. Its slot remains
occupied between lanes until all admitted builds for that PR finish. A
non-PR build, including a schedule on `main`, consumes its own slot and the
same GPU budget. The preserved Slurm dispatcher has its existing independent
limits; this new admission database does not control legacy Slurm jobs.
Lanes retain their existing one-, two-, or four-GPU requirements. Four-GPU
SSIM and training lanes wait for that PR's other lanes to release capacity.
Selected integration lanes wait for the golden gate and are skipped if it
fails. CPU-only GitHub checks are unaffected. Pending admission is bounded by
`queue_timeout_seconds`, initially six hours. The coordinator's Buildkite
command timeout is eight hours, including admission and lane waits after the
command starts; time waiting for a Buildkite agent is outside that timeout.
Reservations persist before a worker is created. Cancellation releases them
only after worker termination is confirmed. A lost coordinator, uncertain
create response, or unavailable backend retains its reservations until
recovery confirms cleanup. There is no heartbeat expiry that silently frees
GPUs while an old worker could still be running.
## Operator Installation
Use an operator-reviewed immutable checkout under
`/opt/fastvideo-gpu-ci/source`. Install the supplied `scripts/gpu_ci/run` and
`scripts/gpu_ci/upload` wrappers as `/opt/fastvideo-gpu-ci/run` and
`/opt/fastvideo-gpu-ci/upload`. They invoke the reviewed Python entrypoint
with isolated import mode. Install the supplied files instead of generating
shell wrappers from build metadata:
```bash
install -m 755 /opt/fastvideo-gpu-ci/source/scripts/gpu_ci/run /opt/fastvideo-gpu-ci/run
install -m 755 /opt/fastvideo-gpu-ci/source/scripts/gpu_ci/upload /opt/fastvideo-gpu-ci/upload
```
The controller needs Python 3.10 or newer, PyYAML,
the Buildkite agent, and `kubectl`; Modal additionally needs the reviewed
Modal SDK and a controller-side Modal credential.
Copy [the configuration example](../../scripts/gpu_ci/config.example.json)
to `/etc/fastvideo-gpu-ci.json`. Keep the installation and configuration
operator-owned and unwritable by workers. Configure the actual legacy
`slurm_uploader` command, replace each enabled backend's image placeholder
with a reviewed registry digest, and leave `default_backend` as `slurm`
during canaries. The checked-in placeholders deliberately cannot start jobs.
The permanent dispatcher must reach the Kubernetes API independently of a
developer laptop, SSH tunnel, or Tailscale session. The existing CPU development
pod is useful for investigating access, but is not a production CI service.
Use a dedicated controller identity and a separate `gpu-ci-dispatch` Buildkite
queue. The trusted bootstrap also needs access to the existing Slurm uploader
while that backend remains enabled.
Run all coordinators on **one host**, with `state_path` and `artifacts_dir` on
its persistent local disk. The implementation uses SQLite transactions and
local file locks. Do not place the database or locks on Lustre/NFS, run
independent database copies, or scale controller replicas across nodes. Those
configurations would invalidate the global limits. Multiple Buildkite agent
processes on the same host can share the installation and ledger; enough
agents are needed for concurrent suite coordinators and queued work.
Configure agent-owned hooks to skip repository checkout and reject commands
other than the trusted uploader and coordinator. Disable repository hooks
and plugins on this queue. Do not use a PR checkout to load `scripts/gpu_ci`,
its configuration, or its worker entrypoint. Build metadata is validated by
the dispatcher; arbitrary environment variables are not forwarded into GPU
workers. Buildkite, Kubernetes, and Modal control-plane credentials stay on
the dispatcher.
For Kubernetes, provision namespace-scoped permissions to create/get/delete
Jobs and read Pods and their logs. Worker Pods disable service-account token
mounting and request the exact GPU count on ARM64 GB200 nodes. The image must
include the expected `/opt/venv` runtime, CUDA/SM100 support, and the reviewed
FA4 dependencies. Use a separate AMD64 image digest for Modal. The Modal
profile preserves H100 for encoder, custom-kernel, and VSA lanes, and L40S
for the remaining lanes. Its reviewed image must support both SM89 and
SM90a. Attention settings are selected per backend and lane to preserve the
existing Modal and GB200 FA4 profiles; do not replace them with one global
attention override. Validate the images against their respective hardware
before enabling merge gates.
Optional Kubernetes configuration includes `context`, `hf_secret`,
`hf_secret_key`, `cache_pvc`, `cache_subpath`, `artifacts_pvc`, and
`artifacts_subpath`. Only provide a read-only Hugging Face credential when
private/gated downloads require it. PR workers are untrusted and can access
any credential supplied to them, so never provide a reference-publication
or other write-capable token. Prefer a pre-populated read-only cache.
Use dedicated CI PVCs and relative CI subpaths. The personal
`lustre-pvc-vllm` and `nfs-pvc-vllm` claims are rejected. The cache is mounted
read-only; mutable references and locks stay inside the worker. The artifact
init container prepares each Job's dedicated artifact subdirectory before
the worker mounts it. An
operator quota or admission policy scoped to CI can provide a second GPU
ceiling; do not apply an eight-GPU quota to the shared `vllm` namespace if it
would also cap other users' work.
## Wire Every Buildkite Entry Pipeline
Set each pipeline's operator-owned bootstrap command to
`/opt/fastvideo-gpu-ci/upload`. Apply this to all three entry pipelines:
| Pipeline | Trigger and scope |
|---|---|
| `pr-fastcheck` | Automatic PR webhook; `TEST_SCOPE=fastcheck` or unset. |
| `ci` | Existing API triggers for Fastcheck, full, merge, direct, and scheduled SSIM. Keep its incoming PR webhook disabled to avoid duplicate builds. |
| `fastvideo-performance-lane` | Existing schedule; `TEST_SCOPE=direct`, `TEST_TYPE=performance`. |
The wrapper resolves `CI_GPU_BACKEND` from the build environment, falling
back to `default_backend` in the trusted configuration. Allowed values are
exactly `slurm`, `modal`, and `vllm`; unknown values fail. For Slurm, upload
delegates to the unchanged private uploader. For Modal or `vllm`, it uploads
one fixed `/opt/fastvideo-gpu-ci/run` command on the dedicated queue, with
the selected backend and scope pinned in step environment.
Set `CI_GPU_BACKEND=vllm` in the Buildkite build environment for a rack canary,
or `modal` for a Modal canary. The existing API build payload can carry the
same string in its `env` object. Automatic PR webhook builds inherit the
operator's configured default unless pipeline/build configuration explicitly
overrides it. Do not put backend selection inside a PR-controlled command.
Keep the existing exact `BUILDKITE_COMMIT`, repository, PR identity, and
`TEST_SCOPE` metadata. Direct runs also need an allowlisted `TEST_TYPE`.
Merge runs need the trusted base-branch planner's `MERGE_TEST_PLAN`,
`MERGE_GOLDEN_TESTS`, and `MERGE_SSIM_TESTS`. Full/direct/scheduled quality
runs keep their complete matrices. The worker fetches and verifies the
immutable commit before installing dependencies or invoking a lane script.
The new Modal adapter uses bounded one-to-four-GPU sandboxes and the same
lane scripts as Kubernetes. It does not reactivate `pr_test.sh` or the old
Modal SSIM fan-out, which cannot enforce this shared four-GPU-per-PR budget.
The legacy files remain available for their existing manual workflows.
## Statuses and Default-Backend Cutover
The selected adapter reports separate suite contexts:
| Scope | GitHub context |
|---|---|
| Fastcheck | `gpu-ci/<backend>/fastcheck-passed` |
| Merge or explicit full suite | `gpu-ci/<backend>/full-suite-passed` |
| Direct lane | `gpu-ci/<backend>/direct-test-completed` |
| Scheduled SSIM | `gpu-ci/<backend>/scheduled-ssim-passed` |
The repository variable `CI_GPU_BACKEND` controls which new backend can
publish the existing required `fastcheck-passed` and `full-suite-passed`
contexts. Empty or `slurm` preserves the existing status behavior. `modal`
or `vllm` enables `ci-gpu-backend-status.yml`, which reads the latest statuses
for that one backend and mirrors both required contexts. A missing result
becomes pending; failure and error remain failures. It never combines one
backend's Fastcheck with another backend's full-suite result. Late status
events re-read current state instead of replaying stale event payloads.
Direct tests are diagnostic and do not promote a whole suite on the new
backends. Rerun the matching Fastcheck or merge/full suite to clear its gate.
Per-build tests on the other new backend remain separate diagnostics. The
legacy direct-test aggregation workflow is disabled while a new backend is
selected. These workflows retain the existing trust assumption that only
authorized status-writing integrations can publish CI status contexts.
For production cutover, drain existing Slurm builds: its preserved pipeline
still emits canonical status contexts. The new uploader rejects a Slurm
override when `default_backend` is `modal` or `vllm`, preventing later legacy
builds from overwriting the promoted backend's checks. Ensure no old bootstrap
bypasses the new uploader. Synchronize the operator `default_backend` and
the GitHub repository variable, then clear or rerun required checks for all
open PRs. Old green contexts do not become new-backend validation merely
because a setting changed. Run both Fastcheck and the merge/full suite on
the promoted backend before allowing merge. Apply the same drain and rerun
procedure when rolling back to Slurm.
## Validation and Recovery
Inspect the rendered opt-in pipeline without submitting a build:
```bash
CI_GPU_BACKEND=vllm TEST_SCOPE=fastcheck \
/opt/fastvideo-gpu-ci/venv/bin/python -I \
/opt/fastvideo-gpu-ci/source/scripts/gpu_ci/entrypoint.py render \
--config /etc/fastvideo-gpu-ci.json
```
Validate a one-GPU lane, then a two-/four-GPU lane and the complete Fastcheck
suite. Exercise two PRs plus a third waiter, concurrent builds of the same
PR, cancellation, retries, and coordinator restart. Confirm observed GPU
reservations never exceed two PRs, four per PR, or eight in total. A GPU
reservation includes a pending worker, so unavailable nodes cannot cause
the controller to submit more work than its allowance.
Run SSIM, training, and performance canaries separately. References must
match the effective GPU/runtime/attention backend; do not silently reuse
L40S performance results as GB200 baselines or reseed references as part of
routine CI. Workers keep W&B offline and disable reference publication.
Inspect the persistent ledger and recover an abandoned attempt with:
```bash
/opt/fastvideo-gpu-ci/venv/bin/python -I \
/opt/fastvideo-gpu-ci/source/scripts/gpu_ci/entrypoint.py status \
--config /etc/fastvideo-gpu-ci.json
/opt/fastvideo-gpu-ci/venv/bin/python -I \
/opt/fastvideo-gpu-ci/source/scripts/gpu_ci/entrypoint.py recover \
--config /etc/fastvideo-gpu-ci.json --build-id BUILD_ID.JOB_ID
```
Use the exact ledger ID from `status`; each Buildkite retry has a distinct
job ID. Recovery refuses a live coordinator, persists cancellation, stops
owned resources, and releases reservations only after confirming termination.
If creation or cleanup is ambiguous, investigate the recorded handle on its
backend and retain the reservation until the outcome is known. Do not delete
the database or manually zero counters to unblock the queue.
The dispatcher uploads controller logs, per-lane numeric results, request
metadata, and the suite summary to Buildkite. These are the sources for its
exit status. Generated videos and JUnit files stay in the worker unless a
dedicated Kubernetes artifact PVC is configured; automatic publication of
those worker files and Modal worker artifacts is not implemented. This is
a deployment limitation to account for before replacing existing artifact
review workflows.
@@ -0,0 +1,317 @@
# SPDX-License-Identifier: Apache-2.0
"""CPU-only admission tests; no FastVideo, GPU, or backend service imports."""
from __future__ import annotations
import importlib.util
import json
import subprocess
import sys
from pathlib import Path
import pytest
ADMISSION_PATH = Path(__file__).resolve().parents[3] / "scripts/gpu_ci/admission.py"
SPEC = importlib.util.spec_from_file_location("gpu_ci_admission", ADMISSION_PATH)
assert SPEC is not None and SPEC.loader is not None
ADMISSION = importlib.util.module_from_spec(SPEC)
SPEC.loader.exec_module(ADMISSION)
AdmissionStore = ADMISSION.AdmissionStore
def _active(store):
return [allocation for allocation in store.snapshot()["allocations"] if allocation["state"] == "active"]
def _register(store, build, pr=None, backend="kubernetes", commit="a" * 40):
store.register(build, pr, backend, commit)
def test_same_pr_builds_share_gpu_budget_and_hold_slot_across_lane_gaps(tmp_path):
store = AdmissionStore(tmp_path / "admission.db")
_register(store, "fastcheck", "repo#1")
_register(store, "other-pr", "repo#2")
assert store.try_admit("fastcheck")
assert store.try_admit("other-pr")
_register(store, "merge", "repo#1", backend="modal")
_register(store, "direct", "repo#1", commit="b" * 40)
_register(store, "waiting", "repo#3")
assert store.try_admit("merge")
assert store.try_admit("direct")
assert not store.try_admit("waiting")
assert store.try_acquire("fastcheck", "encoder", 2, {"job": "encoder"})
assert store.try_acquire("merge", "training", 2, {"call": "training"})
assert not store.try_acquire("direct", "vae", 1, {"job": "vae"})
store.release("fastcheck", "encoder")
store.finish("fastcheck")
store.release("merge", "training")
store.finish("merge")
assert not store.try_admit("waiting")
assert store.try_acquire("direct", "vae", 1, {"job": "vae"})
store.release("direct", "vae")
assert not store.try_admit("waiting")
store.finish("direct")
assert store.try_admit("waiting")
def test_new_pr_admission_is_fifo_even_when_later_pr_polls_first(tmp_path):
store = AdmissionStore(tmp_path / "admission.db")
for index in range(4):
_register(store, f"build-{index}", f"repo#{index}")
assert not store.try_admit("build-2")
assert store.try_admit("build-1")
assert not store.try_admit("build-2")
assert store.try_admit("build-0")
store.finish("build-1")
assert not store.try_admit("build-3")
assert store.try_admit("build-2")
def test_non_pr_builds_get_distinct_slots(tmp_path):
store = AdmissionStore(tmp_path / "admission.db")
_register(store, "schedule-1")
_register(store, "schedule-2", "")
_register(store, "schedule-3")
assert store.try_admit("schedule-1")
assert store.try_admit("schedule-2")
assert not store.try_admit("schedule-3")
store.finish("schedule-1")
assert store.try_admit("schedule-3")
def test_gpu_queue_does_not_starve_large_lanes_or_block_other_pr_budget(tmp_path):
store = AdmissionStore(tmp_path / "admission.db")
_register(store, "a", "repo#1")
_register(store, "b", "repo#2")
assert store.try_admit("a") and store.try_admit("b")
assert store.try_acquire("a", "running", 3, {"job": "running"})
assert not store.try_acquire("a", "older", 2, {"job": "older"})
assert not store.try_acquire("a", "younger", 1, {"job": "younger"})
assert store.try_acquire("b", "independent", 4, {"job": "independent"})
store.release("a", "running")
assert not store.try_acquire("a", "younger", 1, {"job": "younger"})
assert store.try_acquire("a", "older", 2, {"job": "older"})
assert store.try_acquire("a", "younger", 1, {"job": "younger"})
def test_gpu_queue_preserves_global_capacity_for_oldest_eligible_request(tmp_path):
store = AdmissionStore(tmp_path / "admission.db", max_gpus=4)
for build in ("a", "b"):
_register(store, build, f"repo#{build}")
assert store.try_admit(build)
assert store.try_acquire("a", "running", 3, {"job": "running"})
assert not store.try_acquire("b", "older", 2, {"job": "older"})
assert not store.try_acquire("a", "younger", 1, {"job": "younger"})
store.release("a", "running")
assert not store.try_acquire("a", "younger", 1, {"job": "younger"})
assert store.try_acquire("b", "older", 2, {"job": "older"})
assert store.try_acquire("a", "younger", 1, {"job": "younger"})
def test_restart_retains_handles_stale_allocations_and_pr_slots(tmp_path, monkeypatch):
path = tmp_path / "admission.db"
store = AdmissionStore(path, max_prs=1)
with monkeypatch.context() as clock:
clock.setattr(ADMISSION.time, "time", lambda: 1.0)
_register(store, "running", "repo#1")
assert store.try_admit("running")
assert store.try_acquire("running", "lane", 4, {"namespace": "vllm", "job": "known-before-create"})
restarted = AdmissionStore(path, max_prs=1)
_register(restarted, "waiting", "repo#2")
assert not restarted.try_admit("waiting")
assert _active(restarted)[0]["handle"] == {"namespace": "vllm", "job": "known-before-create"}
with pytest.raises(RuntimeError, match="active GPU allocations"):
restarted.finish("running")
restarted.heartbeat("running")
assert len(_active(restarted)) == 1
restarted.release("running", "lane")
assert not restarted.try_admit("waiting")
restarted.finish("running")
assert restarted.try_admit("waiting")
def test_idempotent_reservations_do_not_resurrect_released_work(tmp_path):
store = AdmissionStore(tmp_path / "admission.db")
_register(store, "build", "repo#1")
assert store.try_admit("build")
_register(store, "build", "repo#1")
assert store.try_acquire("build", "lane", 2, {"job": "one", "namespace": "vllm"})
assert store.try_acquire("build", "lane", 2, {"namespace": "vllm", "job": "one"})
assert sum(row["gpus"] for row in _active(store)) == 2
with pytest.raises(ValueError, match="identity changed"):
store.try_acquire("build", "lane", 2, {"job": "two"})
with pytest.raises(ValueError, match="identity changed"):
store.try_acquire("build", "lane", 3, {"job": "one", "namespace": "vllm"})
store.release("build", "lane")
store.release("build", "lane")
assert not store.try_acquire("build", "lane", 2, {"job": "one", "namespace": "vllm"})
store.finish("build")
_register(store, "build", "repo#1")
store.heartbeat("build")
assert not store.try_admit("build")
assert not store.try_acquire("build", "new-lane", 1, {"job": "late"})
assert store.snapshot()["builds"][0]["state"] == "finished"
@pytest.mark.parametrize("changed", [{"pr": "repo#2"}, {"backend": "modal"}, {"commit": "b" * 40}])
def test_existing_build_identity_cannot_be_replaced(tmp_path, changed):
store = AdmissionStore(tmp_path / "admission.db")
_register(store, "build", "repo#1")
values = {"pr": "repo#1", "backend": "kubernetes", "commit": "a" * 40}
values.update(changed)
with pytest.raises(ValueError, match="Build identity changed"):
_register(store, "build", **values)
def test_cancel_blocks_launches_but_keeps_active_reservation_until_confirmed_stop(tmp_path):
path = tmp_path / "admission.db"
store = AdmissionStore(path, max_prs=1)
_register(store, "running", "repo#1")
_register(store, "waiting", "repo#2")
assert store.try_admit("running")
assert store.try_acquire("running", "active", 4, {"job": "active"})
assert not store.try_acquire("running", "queued", 1, {"job": "queued"})
store.request_cancel("running")
restarted = AdmissionStore(path, max_prs=1)
assert restarted.is_cancelled("running")
assert not restarted.try_admit("running")
assert not restarted.try_acquire("running", "active", 4, {"job": "active"})
assert not restarted.try_acquire("running", "new", 1, {"job": "new"})
assert not restarted.try_admit("waiting")
with pytest.raises(RuntimeError, match="active GPU allocations"):
restarted.finish("running")
snapshot = restarted.snapshot()
assert {row["lane_id"]: row["state"] for row in snapshot["allocations"]} == {
"active": "active", "queued": "cancelled"
}
restarted.release("running", "active")
restarted.finish("running")
assert restarted.try_admit("waiting")
def test_cancelled_queued_pr_does_not_block_fifo_and_can_finish(tmp_path):
store = AdmissionStore(tmp_path / "admission.db", max_prs=1)
_register(store, "cancelled", "repo#1")
_register(store, "next", "repo#2")
store.request_cancel("cancelled")
store.request_cancel("cancelled")
assert not store.try_admit("cancelled")
assert store.try_admit("next")
store.finish("cancelled")
store.finish("cancelled")
assert store.snapshot()["builds"][0]["cancelled"]
def test_conflicting_limits_are_rejected_without_changing_ledger(tmp_path):
path = tmp_path / "admission.db"
store = AdmissionStore(path)
_register(store, "build", "repo#1")
for changed in ({"max_prs": 3}, {"max_gpus_per_pr": 5}, {"max_gpus": 9}):
with pytest.raises(ValueError, match="limits conflict"):
AdmissionStore(path, **changed)
assert AdmissionStore(path).snapshot() == store.snapshot()
@pytest.mark.parametrize("gpus", [0, -1, 5, True, 1.5])
def test_impossible_gpu_requests_are_rejected(tmp_path, gpus):
store = AdmissionStore(tmp_path / "admission.db")
_register(store, "build", "repo#1")
assert store.try_admit("build")
with pytest.raises(ValueError, match="GPU"):
store.try_acquire("build", "bad", gpus, {"job": "bad"})
assert not store.snapshot()["allocations"]
def test_unknown_build_operations_fail_closed(tmp_path):
store = AdmissionStore(tmp_path / "admission.db")
for operation in (store.try_admit, store.finish, store.heartbeat, store.request_cancel, store.is_cancelled):
with pytest.raises(KeyError, match="Unknown build"):
operation("unknown")
with pytest.raises(KeyError, match="Unknown build"):
store.try_acquire("unknown", "lane", 1, {"job": "unknown"})
_CONTENDER = """
import importlib.util
import json
import sys
spec = importlib.util.spec_from_file_location('admission', sys.argv[1])
module = importlib.util.module_from_spec(spec)
spec.loader.exec_module(module)
request = json.loads(sys.argv[3])
store = module.AdmissionStore(sys.argv[2], **request['limits'])
print('ready', flush=True)
assert sys.stdin.readline().strip() == 'go'
admitted = store.try_admit(request['build'])
acquired = admitted and store.try_acquire(
request['build'], request['lane'], request['gpus'], {'job': request['handle']})
print(json.dumps({'admitted': admitted, 'acquired': acquired}), flush=True)
"""
def _run_contenders(path, requests, limits):
processes = []
try:
for request in requests:
process = subprocess.Popen(
[sys.executable, "-c", _CONTENDER, str(ADMISSION_PATH), str(path),
json.dumps({**request, "limits": limits})],
stdin=subprocess.PIPE, stdout=subprocess.PIPE, stderr=subprocess.PIPE, text=True,
)
processes.append(process)
for process in processes:
assert process.stdout.readline().strip() == "ready"
for process in processes:
process.stdin.write("go\n")
process.stdin.flush()
results = []
for process in processes:
output, error = process.communicate(timeout=30)
assert process.returncode == 0, error
results.append(json.loads(output))
return results
finally:
for process in processes:
if process.poll() is None:
process.kill()
process.communicate()
@pytest.mark.parametrize("pr_keys,gpus,max_gpus,expected_prs,expected_gpus", [
(["repo#1", "repo#2", "repo#3", "repo#4"], 4, 8, 2, 8),
(["repo#1"] * 8, 1, 8, 1, 4),
(["repo#1"] * 4 + ["repo#2"] * 4, 2, 6, 2, 6),
])
def test_multiple_processes_cannot_exceed_pr_or_gpu_capacity(
tmp_path, pr_keys, gpus, max_gpus, expected_prs, expected_gpus,
):
path = tmp_path / "admission.db"
limits = {"max_prs": 2, "max_gpus_per_pr": 4, "max_gpus": max_gpus}
store = AdmissionStore(path, **limits)
requests = []
for index, pr_key in enumerate(pr_keys):
build = f"build-{index}"
_register(store, build, pr_key, backend="modal" if index % 2 else "kubernetes")
requests.append({"build": build, "lane": "lane", "gpus": gpus, "handle": f"job-{index}"})
_run_contenders(path, requests, limits)
snapshot = store.snapshot()
builds = {build["build_id"]: build for build in snapshot["builds"]}
assert len({build["pr_key"] for build in builds.values() if build["state"] == "admitted"}) == expected_prs
usage = {}
for allocation in _active(store):
pr_key = builds[allocation["build_id"]]["pr_key"]
usage[pr_key] = usage.get(pr_key, 0) + allocation["gpus"]
assert sum(usage.values()) == expected_gpus
assert all(count <= 4 for count in usage.values())
def test_simultaneous_duplicate_retries_share_one_reservation(tmp_path):
path = tmp_path / "admission.db"
store = AdmissionStore(path)
_register(store, "build", "repo#1")
request = {"build": "build", "lane": "lane", "gpus": 4, "handle": "one-worker"}
results = _run_contenders(path, [request] * 6, {})
assert all(result["acquired"] for result in results)
assert len(_active(store)) == 1
assert _active(store)[0]["gpus"] == 4
@@ -0,0 +1,302 @@
"""CPU-only lifecycle and isolation contracts for GPU CI adapters."""
import json
from pathlib import Path
import subprocess
import sys
import threading
from types import SimpleNamespace
import pytest
ROOT = Path(__file__).resolve().parents[3]
sys.path.insert(0, str(ROOT))
from scripts.gpu_ci.backends import KubernetesBackend, ModalBackend, REPOSITORY
IMAGE = "ghcr.io/hao-ai-lab/fastvideo/fastvideo-dev@sha256:" + "a" * 64
REQUEST = {"build_id": "build-one", "repository": REPOSITORY, "commit": "b" * 40,
"pr_number": "12", "scope": "direct", "env": {"LD_PRELOAD": "bad.so"}}
LANE = {"key": "unit", "script": ".buildkite/scripts/unit_test.sh", "gpus": 1,
"extras": ["test"], "fa4": "0", "timeout_seconds": 10}
def kubernetes(tmp_path, **config):
return KubernetesBackend({"image": IMAGE, **config}, ROOT / "scripts/gpu_ci", tmp_path)
def test_manifest_pins_gpu_count_and_does_not_expose_cluster_or_personal_storage(tmp_path):
backend = kubernetes(tmp_path, hf_secret="ci-hf-readonly")
handle = backend.handle("build-one", "ssim")
manifest = backend.manifest(handle, REQUEST, {**LANE, "gpus": 4})
pod = manifest["spec"]["template"]["spec"]
worker = pod["containers"][0]
assert pod["automountServiceAccountToken"] is False
assert not any(pod[key] for key in ("hostNetwork", "hostPID", "hostIPC"))
assert worker["securityContext"]["privileged"] is False
assert worker["resources"]["limits"]["nvidia.com/gpu"] == "4"
assert manifest["spec"]["backoffLimit"] == 0
assert worker["image"] == IMAGE
assert not any("hostPath" in volume or "persistentVolumeClaim" in volume for volume in pod["volumes"])
env = {entry["name"]: entry for entry in worker["env"]}
assert "LD_PRELOAD" not in env
assert env["HF_TOKEN"]["valueFrom"]["secretKeyRef"]["name"] == "ci-hf-readonly"
assert env["FASTVIDEO_CI_LOCAL_ONLY"]["value"] == "1"
assert handle == backend.handle("build-one", "ssim")
@pytest.mark.parametrize("config", [
{"image": "image:latest"}, {"namespace": "default"},
{"cache_pvc": "lustre-pvc-vllm", "cache_subpath": "satyam"},
{"cache_pvc": "ci-cache", "cache_subpath": "../personal"},
{"artifacts_pvc": "ci-output", "artifacts_subpath": "/"},
])
def test_unsafe_backend_configuration_is_rejected(tmp_path, config):
with pytest.raises(ValueError):
kubernetes(tmp_path, **config)
def test_pvc_mounts_are_narrow_and_cache_is_readonly(tmp_path):
backend = kubernetes(tmp_path, cache_pvc="ci-cache", cache_subpath="hf-hub",
artifacts_pvc="ci-output", artifacts_subpath="runs")
handle = backend.handle("build-one", "unit")
pod = backend.manifest(handle, REQUEST, LANE)["spec"]["template"]["spec"]
mounts = {mount["name"]: mount for mount in pod["containers"][0]["volumeMounts"]}
assert mounts["cache"]["readOnly"]
assert mounts["artifacts"]["subPath"] == "runs/" + handle["name"]
init = pod["initContainers"][0]
assert init["command"] == ["mkdir", "-p", "/ci-artifacts/" + handle["name"]]
assert init["volumeMounts"][0]["subPath"] == "runs"
def test_kubectl_create_uses_json_stdin_and_no_shell(tmp_path, monkeypatch):
backend = kubernetes(tmp_path)
calls = []
def run(command, **kwargs):
calls.append((command, kwargs))
return SimpleNamespace(returncode=0, stdout="created", stderr="")
monkeypatch.setattr(subprocess, "run", run)
backend._kubectl(backend.handle("b", "unit"), "create", "-f", "-", body={"safe": True})
command, kwargs = calls[0]
assert command[-3:] == ["create", "-f", "-"]
assert json.loads(kwargs["input"]) == {"safe": True}
assert not kwargs.get("shell", False)
def fake_cluster(backend, handle, monkeypatch, exit_code=7):
state = {"created": False, "deleted": False}
job = {"metadata": {"uid": "job-uid", "labels": {"fastvideo-ci/run-id": handle["run_id"]}}}
pod = {"metadata": {"name": "worker-pod", "ownerReferences": [{"uid": "job-uid"}]},
"status": {"containerStatuses": [{"name": "worker", "state": {"terminated": {"exitCode": exit_code}}}]}}
monkeypatch.setattr(backend, "_job", lambda unused: job if state["created"] else None)
monkeypatch.setattr(backend, "_pods", lambda unused: [pod] if state["created"] else [])
def kubectl(unused, *args, **kwargs):
if args[0] == "create":
state["created"] = True
elif args[0] == "delete":
assert "--cascade=foreground" in args
assert "--wait=true" in args
state.update(created=False, deleted=True)
elif args[0] == "logs":
return "captured worker output\n"
return ""
monkeypatch.setattr(backend, "_kubectl", kubectl)
monkeypatch.setattr(subprocess, "Popen", lambda *args, **kwargs: SimpleNamespace(wait=lambda **kw: 0))
return state, job, pod
@pytest.mark.parametrize("exit_code", [0, 7, 137])
def test_kubernetes_propagates_authoritative_worker_exit(tmp_path, monkeypatch, exit_code):
backend = kubernetes(tmp_path)
handle = backend.handle("b", "unit")
state, _, _ = fake_cluster(backend, handle, monkeypatch, exit_code)
assert backend.run(handle, REQUEST, LANE, threading.Event()) == exit_code
assert backend.stop(handle)
assert state["deleted"]
def test_worker_metadata_uses_validated_request_and_pinned_image(tmp_path):
backend = kubernetes(tmp_path)
request = {**REQUEST, "buildkite_build_id": "real-build", "buildkite_job_id": "real-job",
"env": {"BUILDKITE_REPO": "https://bad.example/repo", "BUILDKITE_COMMIT": "bad",
"BUILDKITE_BUILD_ID": "fake", "BUILDKITE_JOB_ID": "fake",
"BUILDKITE_PULL_REQUEST": "999", "FASTVIDEO_CONTAINER_IMAGE_REF": "image:latest"}}
env = backend._environment(request, LANE)
assert env["BUILDKITE_COMMIT"] == REQUEST["commit"]
assert env["BUILDKITE_REPO"] == REPOSITORY
assert env["BUILDKITE_BUILD_ID"] == "real-build"
assert env["BUILDKITE_JOB_ID"] == "real-job"
assert env["BUILDKITE_PULL_REQUEST"] == "12"
assert env["FASTVIDEO_CONTAINER_IMAGE_REF"] == IMAGE
with pytest.raises(ValueError, match="build/job identity"):
backend._environment({**request, "buildkite_job_id": "bad\nidentifier"}, LANE)
@pytest.mark.parametrize("logs", ["", "error"])
def test_worker_success_without_retrievable_logs_is_infrastructure_failure(tmp_path, monkeypatch, logs):
backend = kubernetes(tmp_path)
handle = backend.handle("b", "unit")
fake_cluster(backend, handle, monkeypatch, 0)
original = backend._kubectl
def kubectl(*args, **kwargs):
if args[1] == "logs":
if logs == "error":
raise RuntimeError("kubectl logs failed")
return logs
return original(*args, **kwargs)
monkeypatch.setattr(backend, "_kubectl", kubectl)
with pytest.raises(RuntimeError, match="logs"):
backend.run(handle, REQUEST, LANE, threading.Event())
assert backend.stop(handle)
def test_final_log_fetch_recovers_failed_stream(tmp_path, monkeypatch):
backend = kubernetes(tmp_path)
handle = backend.handle("b", "unit")
fake_cluster(backend, handle, monkeypatch, 0)
monkeypatch.setattr(subprocess, "Popen", lambda *args, **kwargs: SimpleNamespace(wait=lambda **kw: 1))
assert backend.run(handle, REQUEST, LANE, threading.Event()) == 0
assert "captured worker output" in (tmp_path / f"{handle['name']}.final.log").read_text()
def fake_modal(monkeypatch, backend, handle, exit_code=0):
class NotFoundError(Exception):
pass
state = {"created": None, "terminated": False, "existing": None}
sandbox = SimpleNamespace(
object_id="sb-test", stdout=iter(["test output\n"]), stderr=iter([]),
poll=lambda: exit_code, get_tags=lambda: {"fastvideo-ci-run-id": handle["run_id"]}, detach=lambda: None)
def lookup(*args):
if state["existing"] is None:
raise NotFoundError
return state["existing"]
def create(*args, **kwargs):
state["created"] = (args, kwargs)
state["existing"] = sandbox
return sandbox
def terminate(*, wait):
assert wait is True
state["terminated"] = True
return 137
sandbox.terminate = terminate
image = SimpleNamespace(add_local_file=lambda *args: "worker-image")
module = SimpleNamespace(
exception=SimpleNamespace(NotFoundError=NotFoundError),
Sandbox=SimpleNamespace(from_name=lookup, create=create),
Image=SimpleNamespace(from_registry=lambda image_ref: image),
App=SimpleNamespace(lookup=lambda *args, **kwargs: "app"),
Secret=SimpleNamespace(from_name=lambda name: name))
monkeypatch.setattr(backend, "_modal", lambda: module)
return state
@pytest.mark.parametrize("gpus", [1, 2, 4])
def test_modal_uses_single_bounded_sandbox_and_same_worker(tmp_path, monkeypatch, gpus):
backend = ModalBackend({"image": IMAGE}, ROOT / "scripts/gpu_ci", tmp_path)
handle = backend.handle("b", "unit")
state = fake_modal(monkeypatch, backend, handle, 7)
assert backend.run(handle, REQUEST, {**LANE, "gpus": gpus}, threading.Event()) == 7
args, kwargs = state["created"]
assert args == ("/bin/bash", "/opt/fastvideo-ci/worker.sh")
assert kwargs["gpu"] == f"L40S:{gpus}"
assert kwargs["name"] == handle["name"]
assert kwargs["env"]["FASTVIDEO_CI_COMMIT"] == REQUEST["commit"]
assert backend.stop(handle)
assert state["terminated"]
def test_modal_unknown_existing_sandbox_is_not_reused_or_stopped(tmp_path, monkeypatch):
backend = ModalBackend({"image": IMAGE}, ROOT / "scripts/gpu_ci", tmp_path)
handle = backend.handle("b", "unit")
state = fake_modal(monkeypatch, backend, handle)
state["existing"] = SimpleNamespace(get_tags=lambda: {"fastvideo-ci-run-id": "other"})
with pytest.raises(RuntimeError, match="already exists"):
backend.run(handle, REQUEST, LANE, threading.Event())
assert backend.stop(handle) is False
@pytest.mark.parametrize("key,gpus", [("training-vsa", 2), ("encoder", 1), ("kernel-tests", 1)])
def test_modal_hopper_lanes_keep_reviewed_hardware(tmp_path, monkeypatch, key, gpus):
backend = ModalBackend({"image": IMAGE}, ROOT / "scripts/gpu_ci", tmp_path)
handle = backend.handle("b", key)
state = fake_modal(monkeypatch, backend, handle)
lane = {**LANE, "key": key, "gpus": gpus, "modal_gpu": "H100", "modal_fa4": "0"}
assert backend.run(handle, REQUEST, lane, threading.Event()) == 0
assert state["created"][1]["gpu"] == f"H100:{gpus}"
assert backend.stop(handle)
def test_modal_attention_policy_does_not_mutate_shared_gb200_lane(tmp_path):
modal = ModalBackend({"image": IMAGE}, ROOT / "scripts/gpu_ci", tmp_path / "modal")
kube = kubernetes(tmp_path / "kubernetes")
lane = {**LANE, "modal_gpu": "L40S", "modal_fa4": "1"}
assert modal._environment(REQUEST, lane)["FASTVIDEO_FA4"] == "1"
assert kube._environment(REQUEST, lane)["FASTVIDEO_FA4"] == "0"
assert lane["fa4"] == "0"
assert modal._environment(REQUEST, lane)["FASTVIDEO_CONTAINER_IMAGE_REF"] == IMAGE
assert modal.image == IMAGE
@pytest.mark.parametrize("override", [{"modal_gpu": "A100"}, {"modal_gpu": "H100:8"}, {"modal_fa4": "2"}])
def test_modal_rejects_unreviewed_hardware_or_attention_policy(tmp_path, override):
backend = ModalBackend({"image": IMAGE}, ROOT / "scripts/gpu_ci", tmp_path)
with pytest.raises(ValueError):
backend._environment(REQUEST, {**LANE, **override})
def test_worker_rejects_repository_override_before_checkout(tmp_path):
result = subprocess.run(["bash", str(ROOT / "scripts/gpu_ci/worker.sh")],
env={"PATH": "/usr/bin:/bin", "FASTVIDEO_CI_REPOSITORY": "https://example.invalid/repo"},
cwd=tmp_path, capture_output=True)
assert result.returncode != 0
assert not list(tmp_path.iterdir())
def test_missing_exit_status_is_not_a_pass(tmp_path, monkeypatch):
backend = kubernetes(tmp_path)
handle = backend.handle("b", "unit")
fake_cluster(backend, handle, monkeypatch, None)
with pytest.raises(RuntimeError, match="numeric exit"):
backend.run(handle, REQUEST, LANE, threading.Event())
def test_unknown_job_and_orphaned_pods_do_not_release_lease(tmp_path, monkeypatch):
backend = kubernetes(tmp_path)
handle = backend.handle("b", "unit")
state, job, pod = fake_cluster(backend, handle, monkeypatch)
state["created"] = True
job["metadata"]["labels"]["fastvideo-ci/run-id"] = "someone-else"
assert backend.stop(handle) is False
assert not state["deleted"]
monkeypatch.setattr(backend, "_job", lambda unused: None)
monkeypatch.setattr(backend, "_pods", lambda unused: [pod])
assert backend.stop(handle) is False
def test_uncertain_create_response_remains_reserved(tmp_path, monkeypatch):
backend = kubernetes(tmp_path)
handle = backend.handle("b", "unit")
monkeypatch.setattr(backend, "_job", lambda unused: None)
monkeypatch.setattr(backend, "_pods", lambda unused: [])
(tmp_path / f"{handle['name']}.creating").touch()
assert backend.stop(handle) is False
def test_kubernetes_cancellation_confirms_deletion(tmp_path, monkeypatch):
backend = kubernetes(tmp_path)
handle = backend.handle("b", "unit")
state, _, _ = fake_cluster(backend, handle, monkeypatch)
answers = iter([False, True, True])
event = SimpleNamespace(is_set=lambda: next(answers))
assert backend.run(handle, REQUEST, LANE, event) == 130
assert state["deleted"]
@@ -0,0 +1,323 @@
# SPDX-License-Identifier: Apache-2.0
"""Suite lifecycle coverage with real local admission/locks and fake workers."""
from __future__ import annotations
import json
import sys
import threading
import time
from concurrent.futures import ThreadPoolExecutor
from pathlib import Path
import pytest
REPO_ROOT = Path(__file__).resolve().parents[3]
if str(REPO_ROOT) not in sys.path:
sys.path.insert(0, str(REPO_ROOT))
from scripts.gpu_ci import dispatcher
from scripts.gpu_ci.policy import LANES
class FakeBackend:
"""Check reservation ordering at the external worker creation boundary."""
def __init__(self, store, *, outcomes=None, stop_result=True, hold=None):
self.store = store
self.outcomes = outcomes or {}
self.stop_result = stop_result
self.hold = hold
self.started = threading.Event()
self.events = []
self.persisted_handles = []
self.lock = threading.Lock()
def handle(self, build_id, lane_id):
return {"build": build_id, "lane": lane_id, "job": f"{build_id}-{lane_id}"}
def run(self, handle, request, lane, cancel):
active = [allocation for allocation in self.store.snapshot()["allocations"]
if allocation["build_id"] == request["build_id"]
and allocation["lane_id"] == lane["key"] and allocation["state"] == "active"]
assert len(active) == 1 and active[0]["handle"] == handle
with self.lock:
self.events.append(("start", lane["key"]))
self.persisted_handles.append(handle)
self.started.set()
if self.hold is not None:
deadline = time.monotonic() + 5
while not self.hold.wait(0.005):
if cancel.is_set():
return 143
if time.monotonic() > deadline:
raise RuntimeError("Test worker was not released")
outcome = self.outcomes.get(lane["key"], 0)
if isinstance(outcome, Exception):
raise outcome
return outcome
def stop(self, handle):
with self.lock:
self.events.append(("stop", handle["lane"]))
if isinstance(self.stop_result, Exception):
raise self.stop_result
return self.stop_result
@pytest.fixture
def config(tmp_path):
return {"state_path": str(tmp_path / "admission.db"), "artifacts_dir": str(tmp_path / "artifacts"),
"queue_timeout_seconds": 10, "max_active_prs": 2, "max_gpus_per_pr": 4, "max_gpus": 8}
def _lane(key="unit", gpus=1, fastcheck=True):
return {"key": key, "gpus": gpus, "fastcheck": fastcheck}
def _request(build="build", pr="repo#1", lanes=None, commit="a" * 40):
return {"build_id": build, "pr_key": pr, "commit": commit, "backend": "vllm",
"scope": "full", "lanes": [_lane()] if lanes is None else lanes}
def _summary(config, build="build"):
return json.loads((dispatcher.run_directory(config, build) / "summary.json").read_text())
def _build(store, build="build"):
return next(item for item in store.snapshot()["builds"] if item["build_id"] == build)
def _active(store, build="build"):
return [allocation for allocation in store.snapshot()["allocations"]
if allocation["build_id"] == build and allocation["state"] == "active"]
def _wait_for(predicate):
deadline = time.monotonic() + 5
while not predicate():
if time.monotonic() > deadline:
raise AssertionError("Timed out waiting for dispatcher state")
time.sleep(0.005)
def test_full_suite_success_persists_handles_and_finishes_every_lane(config):
store = dispatcher.make_store(config)
backend = FakeBackend(store)
assert dispatcher.execute(_request(lanes=[dict(lane) for lane in LANES]), config,
backend=backend, poll_seconds=0.001) == 0
summary = _summary(config)
assert set(summary["lanes"]) == {lane["key"] for lane in LANES}
assert all(result["state"] == "passed" and result["stopped"] for result in summary["lanes"].values())
assert len(backend.persisted_handles) == 20
assert _build(store)["state"] == "finished"
assert not _active(store)
assert not summary["recovery_required"]
assert backend.events.index(("stop", "golden-gate")) < backend.events.index(("start", "ssim"))
@pytest.mark.parametrize("outcome, expected", [
(42, 42), (RuntimeError("worker transport failed"), 97), (False, 97), (-1, 97), (256, 97),
])
def test_lane_failure_returns_exit_status_after_confirmed_cleanup(config, outcome, expected):
store = dispatcher.make_store(config)
backend = FakeBackend(store, outcomes={"unit": outcome})
assert dispatcher.execute(_request(), config, backend=backend, poll_seconds=0.001) == expected
assert _summary(config)["lanes"]["unit"]["state"] == "failed"
assert _summary(config)["lanes"]["unit"]["exit_code"] == expected
assert _build(store)["state"] == "finished"
assert not _active(store)
assert ("stop", "unit") in backend.events
def test_failed_golden_gate_skips_integration_but_runs_fastcheck(config):
store = dispatcher.make_store(config)
backend = FakeBackend(store, outcomes={"golden-gate": 7})
lanes = [_lane("ssim", 4, False), _lane("encoder"), _lane("golden-gate", 1, False),
_lane("training", 4, False)]
assert dispatcher.execute(_request(lanes=lanes), config, backend=backend, poll_seconds=0.001) == 7
started = {lane for action, lane in backend.events if action == "start"}
assert started == {"encoder", "golden-gate"}
summary = _summary(config)
assert summary["lanes"]["ssim"]["state"] == "skipped"
assert summary["lanes"]["training"]["reason"] == "golden gate failed"
assert _build(store)["state"] == "finished"
def test_direct_integration_lane_without_selected_golden_runs(config):
store = dispatcher.make_store(config)
backend = FakeBackend(store)
request = _request(lanes=[_lane("ssim", 4, False)])
request["scope"] = "direct"
assert dispatcher.execute(request, config, backend=backend, poll_seconds=0.001) == 0
assert ("start", "ssim") in backend.events
def test_empty_merge_plan_does_not_wait_for_gpu_admission(config):
config["max_active_prs"] = 1
config["queue_timeout_seconds"] = 0
store = dispatcher.make_store(config)
store.register("busy", "repo#other", "vllm", "b" * 40)
assert store.try_admit("busy")
backend = FakeBackend(store)
request = _request(lanes=[])
request["scope"] = "merge"
assert dispatcher.execute(request, config, backend=backend, poll_seconds=0.001) == 0
assert not backend.events
assert not _summary(config)["recovery_required"]
assert _summary(config)["lanes"] == {}
assert _build(store)["state"] == "finished"
def test_persisted_cancellation_of_queued_build_never_starts_backend(config):
config["max_active_prs"] = 1
store = dispatcher.make_store(config)
store.register("busy", "repo#other", "vllm", "b" * 40)
assert store.try_admit("busy")
backend = FakeBackend(store)
cancel = threading.Event()
with ThreadPoolExecutor(max_workers=1) as pool:
future = pool.submit(dispatcher.execute, _request(), config, backend=backend,
cancel=cancel, poll_seconds=0.001)
try:
_wait_for(lambda: any(build["build_id"] == "build" for build in store.snapshot()["builds"]))
store.request_cancel("build")
assert future.result(timeout=5) == dispatcher.INFRASTRUCTURE_FAILURE
finally:
cancel.set()
assert not backend.events
assert _build(store)["state"] == "finished"
assert _build(store)["cancelled"]
assert _build(store, "busy")["state"] == "admitted"
def test_queue_deadline_finishes_without_allocating_gpus(config):
config["queue_timeout_seconds"] = 0
store = dispatcher.make_store(config)
backend = FakeBackend(store)
assert dispatcher.execute(_request(), config, backend=backend, poll_seconds=0.001) == 97
assert not backend.events
assert _build(store)["state"] == "finished"
assert _build(store)["cancelled"]
def test_cancellation_stops_active_worker_and_cancels_pending_lanes(config):
store = dispatcher.make_store(config)
hold = threading.Event()
cancel = threading.Event()
backend = FakeBackend(store, hold=hold)
request = _request(lanes=[_lane("running", 4), _lane("pending", 1)])
with ThreadPoolExecutor(max_workers=1) as pool:
future = pool.submit(dispatcher.execute, request, config, backend=backend,
cancel=cancel, poll_seconds=0.001)
try:
assert backend.started.wait(5)
assert len(_active(store)) == 1
store.request_cancel("build")
assert future.result(timeout=5) != 0
finally:
cancel.set()
hold.set()
assert {lane for action, lane in backend.events if action == "start"} == {"running"}
assert _summary(config)["lanes"]["pending"]["state"] == "cancelled"
assert not _active(store)
assert _build(store)["state"] == "finished"
@pytest.mark.parametrize("stop_result", [False, RuntimeError("termination probe unavailable")])
def test_uncertain_cleanup_keeps_reservation_until_explicit_recovery(config, stop_result):
store = dispatcher.make_store(config)
backend = FakeBackend(store, stop_result=stop_result)
assert dispatcher.execute(_request(), config, backend=backend, poll_seconds=0.001) == 97
assert len(_active(store)) == 1
assert _build(store)["state"] == "admitted"
assert _summary(config)["recovery_required"]
with pytest.raises(RuntimeError, match="termination"):
dispatcher.recover(config, "build", backend=backend)
assert len(_active(store)) == 1
assert store.is_cancelled("build")
backend.stop_result = True
dispatcher.recover(config, "build", backend=backend)
assert not _active(store)
assert _build(store)["state"] == "finished"
dispatcher.recover(config, "build", backend=backend)
def test_recovery_refuses_a_live_dispatcher_and_keeps_its_reservation(config):
store = dispatcher.make_store(config)
store.register("build", "repo#1", "vllm", "a" * 40)
assert store.try_admit("build")
backend = FakeBackend(store)
handle = backend.handle("build", "unit")
assert store.try_acquire("build", "unit", 1, handle)
with dispatcher.build_lock(dispatcher.run_directory(config, "build")):
with pytest.raises(RuntimeError, match="dispatcher is still running"):
dispatcher.recover(config, "build", backend=backend)
assert not backend.events
assert len(_active(store)) == 1
dispatcher.recover(config, "build", backend=backend)
assert not _active(store)
def test_finished_attempt_cannot_be_replayed(config):
store = dispatcher.make_store(config)
backend = FakeBackend(store)
assert dispatcher.execute(_request(), config, backend=backend, poll_seconds=0.001) == 0
before = list(backend.events)
with pytest.raises(RuntimeError, match="already exists"):
dispatcher.execute(_request(), config, backend=backend, poll_seconds=0.001)
assert backend.events == before
def test_overlapping_same_pr_builds_share_slot_and_gpu_budget(config):
config["max_active_prs"] = 1
store = dispatcher.make_store(config)
hold = threading.Event()
cancel = threading.Event()
first_backend = FakeBackend(store, hold=hold)
second_backend = FakeBackend(store, hold=hold)
first = _request("fastcheck", lanes=[_lane("encoder", 2)])
second = _request("merge", lanes=[_lane("training", 2, False)])
second["backend"] = "modal"
with ThreadPoolExecutor(max_workers=2) as pool:
first_future = pool.submit(dispatcher.execute, first, config, backend=first_backend,
cancel=cancel, poll_seconds=0.001)
second_future = None
try:
assert first_backend.started.wait(5)
second_future = pool.submit(dispatcher.execute, second, config, backend=second_backend,
cancel=cancel, poll_seconds=0.001)
assert second_backend.started.wait(5)
store.register("other-pr", "repo#other", "vllm", "b" * 40)
assert not store.try_admit("other-pr")
assert sum(allocation["gpus"] for allocation in _active(store, "fastcheck")
+ _active(store, "merge")) == 4
assert _build(store, "fastcheck")["state"] == "admitted"
assert _build(store, "merge")["state"] == "admitted"
hold.set()
assert first_future.result(timeout=5) == 0
assert second_future.result(timeout=5) == 0
finally:
hold.set()
cancel.set()
assert store.try_admit("other-pr")
def test_unverified_different_sha_does_not_cancel_existing_pr_build(config):
# Arrival order does not establish which commit is the current PR head:
# a delayed webhook or explicit old-SHA diagnostic may arrive last.
store = dispatcher.make_store(config)
store.register("current-head", "repo#1", "vllm", "b" * 40)
assert store.try_admit("current-head")
backend = FakeBackend(store)
delayed = _request("delayed-older-sha", commit="a" * 40)
assert dispatcher.execute(delayed, config, backend=backend, poll_seconds=0.001) == 0
assert not store.is_cancelled("current-head")
assert _build(store, "current-head")["state"] == "admitted"
def test_recovery_of_unknown_attempt_fails_without_backend_actions(config):
backend = FakeBackend(dispatcher.make_store(config))
with pytest.raises(ValueError, match="Unknown build"):
dispatcher.recover(config, "not-registered", backend=backend)
assert not backend.events
@@ -0,0 +1,152 @@
# SPDX-License-Identifier: Apache-2.0
"""CPU-only policy coverage for selectable Buildkite GPU backends."""
from __future__ import annotations
import importlib.util
import json
import shutil
import subprocess
from pathlib import Path
import pytest
import yaml
REPO_ROOT = Path(__file__).resolve().parents[3]
SPEC = importlib.util.spec_from_file_location("gpu_ci_pipeline", REPO_ROOT / "scripts/gpu_ci/pipeline.py")
assert SPEC is not None and SPEC.loader is not None
PIPELINE = importlib.util.module_from_spec(SPEC)
SPEC.loader.exec_module(PIPELINE)
@pytest.mark.parametrize("backend", ["", "slurm"])
@pytest.mark.parametrize("scope", ["", "fastcheck", "full", "merge", "direct", "scheduled"])
def test_default_preserves_complete_existing_pipeline(backend, scope):
canonical = yaml.safe_load((REPO_ROOT / ".buildkite/pipeline.yml").read_text())
rendered = PIPELINE.render_pipeline(canonical, backend, scope)
assert rendered == canonical
assert rendered is not canonical
rendered["steps"][0]["env"]["TEST_TYPE"] = "changed"
assert canonical["steps"][0]["env"]["TEST_TYPE"] != "changed"
@pytest.mark.parametrize("backend", ["modal", "vllm"])
@pytest.mark.parametrize("scope, context", [
("", "fastcheck-passed"),
("fastcheck", "fastcheck-passed"),
("full", "full-suite-passed"),
("merge", "full-suite-passed"),
("direct", "direct-test-completed"),
("scheduled", "scheduled-ssim-passed"),
])
def test_backend_uses_single_trusted_suite_step_and_distinct_status(backend, scope, context):
# A PR can alter its pipeline file, but none of these fields may escape
# into the operator-owned dispatcher pipeline.
canonical = {
"env": {"BASH_ENV": "/checkout/attack.sh", "CI_GPU_BACKEND": "slurm"},
"steps": [{"command": "arbitrary PR command", "plugins": ["unsafe"]}],
"notify": [{"github_commit_status": {"context": "full-suite-passed"}}],
}
rendered = PIPELINE.render_pipeline(canonical, backend, scope)
assert rendered["notify"] == [{"github_commit_status": {"context": f"gpu-ci/{backend}/{context}"}}]
assert set(rendered) == {"notify", "steps"}
assert len(rendered["steps"]) == 1
step = rendered["steps"][0]
assert step["command"] == "/opt/fastvideo-gpu-ci/run"
assert step["agents"] == {"queue": "gpu-ci-dispatch"}
assert step["env"] == {"CI_GPU_BACKEND": backend, "TEST_SCOPE": scope or "fastcheck"}
assert step["timeout_in_minutes"] == 480
assert set(step) == {"label", "key", "command", "timeout_in_minutes", "env", "agents"}
assert ":microscope:" not in step["label"]
assert ":test_tube:" not in step["label"]
assert ":bar_chart:" not in step["label"]
@pytest.mark.parametrize("backend", ["k8s", "VLLM", "vllm; touch /tmp/injected", "$(env)"])
def test_unknown_or_command_like_backend_fails_closed(backend):
with pytest.raises(ValueError, match="backend"):
PIPELINE.render_pipeline({}, backend, "fastcheck")
@pytest.mark.parametrize("scope", ["all", "full; true", "$(env)", "FULL"])
def test_unknown_or_command_like_scope_fails_closed(scope):
with pytest.raises(ValueError, match="scope"):
PIPELINE.render_pipeline({}, "vllm", scope)
@pytest.mark.parametrize("queue", ["", "queue with spaces", "${QUEUE}", "../ci", "q" * 65])
def test_invalid_dispatch_queue_fails_closed(queue):
with pytest.raises(ValueError, match="queue"):
PIPELINE.render_pipeline({}, "vllm", "fastcheck", queue)
def test_trusted_operator_can_choose_a_dedicated_queue():
rendered = PIPELINE.render_pipeline({}, "vllm", "fastcheck", "gpu-ci-dispatch-canary")
assert rendered["steps"][0]["agents"] == {"queue": "gpu-ci-dispatch-canary"}
def _promoted_statuses(statuses, backend="vllm"):
node = shutil.which("node")
if node is None:
pytest.skip("Node.js is required to execute the GitHub status workflow's JavaScript")
workflow = yaml.safe_load((REPO_ROOT / ".github/workflows/ci-gpu-backend-status.yml").read_text())
script = workflow["jobs"]["promote"]["steps"][0]["with"]["script"]
harness = """
const fs = require('node:fs');
const input = JSON.parse(fs.readFileSync(0, 'utf8'));
const updates = [];
process.env.SELECTED_BACKEND = input.backend;
const github = {
paginate: async () => input.statuses,
rest: {repos: {
listCommitStatusesForRef: () => {},
createCommitStatus: async update => {updates.push(update);},
}},
};
const context = {repo: {owner: 'hao-ai-lab', repo: 'FastVideo'}, payload: {sha: 'a'.repeat(40)}};
const AsyncFunction = Object.getPrototypeOf(async function() {}).constructor;
new AsyncFunction('github', 'context', input.script)(github, context)
.then(() => process.stdout.write(JSON.stringify(updates)))
.catch(error => {process.stderr.write(String(error)); process.exitCode = 1;});
"""
result = subprocess.run([node, "-e", harness], text=True, capture_output=True, check=True,
input=json.dumps({"script": script, "statuses": statuses, "backend": backend}))
return {status["context"]: status["state"] for status in json.loads(result.stdout)}
def _status(context, state="success", second=0, identifier=1):
return {"context": context, "state": state, "updated_at": f"2026-09-30T01:00:{second:02d}Z", "id": identifier}
def test_backend_statuses_cannot_combine_to_pass_merge_gate():
statuses = [
_status("gpu-ci/vllm/fastcheck-passed"),
_status("gpu-ci/modal/full-suite-passed"),
_status("full-suite-passed"),
_status("gpu-ci/vllm/direct-test-completed"),
]
assert _promoted_statuses(statuses) == {"fastcheck-passed": "success", "full-suite-passed": "pending"}
@pytest.mark.parametrize("state", ["pending", "failure", "error"])
def test_newer_selected_backend_non_success_replaces_older_success(state):
statuses = [
_status("gpu-ci/vllm/fastcheck-passed", second=1),
_status("gpu-ci/vllm/full-suite-passed", second=1),
_status("gpu-ci/vllm/full-suite-passed", state, second=2, identifier=2),
]
assert _promoted_statuses(statuses) == {"fastcheck-passed": "success", "full-suite-passed": state}
def test_status_id_breaks_timestamp_ties_for_same_commit():
statuses = [
_status("gpu-ci/modal/fastcheck-passed", "success", identifier=1),
_status("gpu-ci/modal/fastcheck-passed", "failure", identifier=2),
_status("gpu-ci/modal/full-suite-passed", "success", identifier=3),
]
assert _promoted_statuses(statuses, "modal") == {"fastcheck-passed": "failure", "full-suite-passed": "success"}
def test_missing_selected_backend_never_reuses_canonical_success():
statuses = [_status("fastcheck-passed"), _status("full-suite-passed")]
assert _promoted_statuses(statuses) == {"fastcheck-passed": "pending", "full-suite-passed": "pending"}
@@ -0,0 +1,115 @@
"""CPU-only request, deployment-config, and trusted CLI boundary tests."""
import json
from pathlib import Path
from unittest.mock import patch
import pytest
import yaml
from scripts.gpu_ci import __main__ as cli
from scripts.gpu_ci.policy import DEFAULT_REPOSITORY, LANES, build_request
ROOT = Path(__file__).resolve().parents[3]
def environment(**overrides):
return {"CI_GPU_BACKEND": "vllm", "BUILDKITE_REPO": DEFAULT_REPOSITORY,
"BUILDKITE_COMMIT": "a" * 40, "BUILDKITE_BUILD_ID": "build-1",
"BUILDKITE_JOB_ID": "job-1", "BUILDKITE_PULL_REQUEST": "12", **overrides}
def test_lane_policy_matches_existing_complete_graph():
pipeline = yaml.safe_load((ROOT / ".buildkite/pipeline.yml").read_text())
assert {lane["key"] for lane in LANES} == {step["key"] for step in pipeline["steps"]}
assert sum(lane["fastcheck"] for lane in LANES) == 6
assert all((ROOT / ".buildkite/scripts" / lane["script"]).is_file() for lane in LANES)
@pytest.mark.parametrize("scope,count", [("fastcheck", 6), ("full", 20), ("scheduled", 1)])
def test_scope_selects_existing_payloads(scope, count):
request = build_request(environment(TEST_SCOPE=scope), {}, ROOT)
assert len(request["lanes"]) == count
assert all(lane["script"].startswith(".buildkite/scripts/") for lane in request["lanes"])
assert request["pr_key"] == DEFAULT_REPOSITORY + "#12"
@pytest.mark.parametrize("lane", LANES, ids=lambda lane: lane["key"])
def test_direct_public_and_compatibility_names(lane):
for name in (lane["public_type"], lane["public_type"] + "_ci", lane["key"]):
request = build_request(environment(TEST_SCOPE="direct", TEST_TYPE=name), {}, ROOT)
assert [item["key"] for item in request["lanes"]] == [lane["key"]]
def test_docs_only_merge_plan_is_explicitly_empty():
assert build_request(environment(TEST_SCOPE="merge", MERGE_TEST_PLAN=",none,"), {}, ROOT)["lanes"] == []
@pytest.mark.parametrize("plan", ["", ",", ",,", ",none,,", ",,golden-gate,,", ",ssim,ssim,",
",unit,", ",none,ssim,", ",unknown,", "ssim", ",../ssim,"])
def test_malformed_merge_plan_fails_closed(plan):
with pytest.raises(ValueError):
build_request(environment(TEST_SCOPE="merge", MERGE_TEST_PLAN=plan), {}, ROOT)
def test_focused_new_pr_test_can_be_validated_inside_worker():
request = build_request(environment(TEST_SCOPE="merge", MERGE_TEST_PLAN=",golden-gate,ssim,",
MERGE_GOLDEN_TESTS="test_new_model.py", MERGE_SSIM_TESTS="all"), {}, ROOT)
assert request["env"]["FASTVIDEO_GOLDEN_TEST_FILES"] == "test_new_model.py"
assert request["env"]["FASTVIDEO_SSIM_TEST_FILES"] == "all"
@pytest.mark.parametrize("files", ["", "../test_model.py", "test_a.py,test_a.py", "$(command)", "test_x.py/evil"])
def test_quality_selection_rejects_missing_or_unsafe_basename(files):
with pytest.raises(ValueError):
build_request(environment(TEST_SCOPE="merge", MERGE_TEST_PLAN=",ssim,", MERGE_SSIM_TESTS=files), {}, ROOT)
@pytest.mark.parametrize("field,value", [("BUILDKITE_REPO", "https://github.com/other/repo.git"),
("BUILDKITE_COMMIT", "main"), ("BUILDKITE_JOB_ID", "../escape"),
("BUILDKITE_PULL_REQUEST", "-1"), ("TEST_SCOPE", "unknown"),
("CI_GPU_BACKEND", "vllm;cmd"), ("FASTVIDEO_SSIM_BOOTSTRAP_MODE", "1")])
def test_untrusted_request_fields_are_rejected(field, value):
with pytest.raises(ValueError):
build_request(environment(**{field: value}), {}, ROOT)
def test_job_retry_has_distinct_attempt_but_shares_pr():
first = build_request(environment(), {}, ROOT)
retry = build_request(environment(BUILDKITE_JOB_ID="job-2"), {}, ROOT)
assert first["build_id"] != retry["build_id"]
assert first["pr_key"] == retry["pr_key"]
def test_credentials_and_shell_controls_never_enter_worker_request():
request = build_request(environment(BUILDKITE_AGENT_TOKEN="secret", HF_TOKEN="secret", BASH_ENV="bad"), {}, ROOT)
assert "secret" not in json.dumps(request)
assert "BASH_ENV" not in request["env"]
def test_config_does_not_allow_raising_agreed_limits(tmp_path):
config = {"state_path": str(tmp_path / "state"), "artifacts_dir": str(tmp_path / "runs")}
path = tmp_path / "config.json"
for key, value in (("max_active_prs", 3), ("max_gpus_per_pr", 5), ("max_gpus", 9)):
path.write_text(json.dumps({**config, key: value}))
with pytest.raises(ValueError, match="ceiling"):
cli.load_config(path)
def test_promoted_backend_blocks_legacy_status_override(tmp_path):
path = tmp_path / "config.json"
path.write_text(json.dumps({"default_backend": "vllm", "state_path": str(tmp_path / "state"),
"artifacts_dir": str(tmp_path / "runs"), "slurm_uploader": ["/trusted/upload"]}))
with patch.dict(cli.os.environ, {"CI_GPU_BACKEND": "slurm"}, clear=True), patch.object(cli.subprocess, "run") as run:
assert cli.main(["upload", "--config", str(path)]) == 97
run.assert_not_called()
def test_original_uploader_is_delegated_without_checkout_or_shell(tmp_path):
path = tmp_path / "config.json"
path.write_text(json.dumps({"state_path": str(tmp_path / "state"), "artifacts_dir": str(tmp_path / "runs"),
"slurm_uploader": ["/trusted/upload"]}))
with patch.dict(cli.os.environ, {}, clear=True), patch.object(cli.subprocess, "run") as run:
run.return_value.returncode = 0
assert cli.main(["upload", "--config", str(path)]) == 0
run.assert_called_once_with(["/trusted/upload"], check=False)
@@ -147,6 +147,10 @@ def _detect_arch_from_torch() -> str:
major, minor = torch.cuda.get_device_capability(0)
if major == 9 and minor == 0:
return "9.0a"
if major == 10 and minor in (0, 3):
# Match build.sh: data-center Blackwell VSA kernels require the
# architecture-specific target, not the generic compute capability.
return f"{major}.{minor}a"
if major == 12 and minor == 0:
return "12.0a"
return f"{major}.{minor}"
@@ -5,9 +5,12 @@ from __future__ import annotations
import errno
import importlib.util
import json
import shutil
import subprocess
import sys
import zipfile
from pathlib import Path
from types import SimpleNamespace
import pytest
import yaml
@@ -168,7 +171,7 @@ def test_ci_runner_image_targets_arm64_sm100_with_opencv_runtime() -> None:
assert job["with"]["tag_suffix"] == "py3.12-cuda13.0.0-sm100"
assert "CUDA_VERSION=13.0.0" in job["with"]["build_args"]
assert "UV_TORCH_BACKEND=cu130" in job["with"]["build_args"]
assert "TORCH_CUDA_ARCH_LIST=10.0" in job["with"]["build_args"]
assert "TORCH_CUDA_ARCH_LIST=10.0a" in job["with"]["build_args"].splitlines()
dockerfile = (REPO_ROOT / "docker/Dockerfile").read_text(encoding="utf-8")
assert " ffmpeg \\\n" in dockerfile
@@ -191,23 +194,68 @@ def test_docker_image_bakes_modal_apt_and_rust_layer() -> None:
assert "ENV PATH=/root/.cargo/bin:${PATH}" in dockerfile
def test_cache_key_uses_resolved_arch_not_raw_env(monkeypatch, tmp_path) -> None:
_patch_stable_metadata(monkeypatch)
monkeypatch.setattr(kernel_build_cache, "_detect_arch_from_torch", lambda: "9.0a")
def _patch_cuda_capability(monkeypatch, capability) -> None:
cuda = SimpleNamespace(is_available=lambda: True, get_device_capability=lambda device: capability)
monkeypatch.setitem(sys.modules, "torch", SimpleNamespace(cuda=cuda))
monkeypatch.setenv("TORCH_CUDA_ARCH_LIST", "9.0a")
explicit_hopper = kernel_build_cache._build_metadata(tmp_path)
@pytest.mark.parametrize("capability, expected", [
((8, 9), "8.9"),
((9, 0), "9.0a"),
((10, 0), "10.0a"),
((10, 3), "10.3a"),
((12, 0), "12.0a"),
((12, 1), "12.1"),
])
def test_detected_arch_enables_device_specific_kernels(monkeypatch, capability, expected) -> None:
_patch_cuda_capability(monkeypatch, capability)
assert kernel_build_cache._detect_arch_from_torch() == expected
@pytest.mark.parametrize("arch", ["9.0a", "10.0a", "10.3a"])
def test_cache_key_uses_resolved_arch_not_raw_env(monkeypatch, tmp_path, arch) -> None:
_patch_stable_metadata(monkeypatch)
# An explicit target must win even when the detected GPU differs.
monkeypatch.setattr(kernel_build_cache, "_detect_arch_from_torch", lambda: "8.9")
monkeypatch.setenv("TORCH_CUDA_ARCH_LIST", arch)
explicit = kernel_build_cache._build_metadata(tmp_path)
monkeypatch.delenv("TORCH_CUDA_ARCH_LIST", raising=False)
detected_hopper = kernel_build_cache._build_metadata(tmp_path)
monkeypatch.setattr(kernel_build_cache, "_detect_arch_from_torch", lambda: arch)
detected = kernel_build_cache._build_metadata(tmp_path)
monkeypatch.setattr(kernel_build_cache, "_detect_arch_from_torch", lambda: "8.9")
detected_l40s = kernel_build_cache._build_metadata(tmp_path)
monkeypatch.setattr(kernel_build_cache, "_detect_arch_from_torch", lambda: arch.removesuffix("a"))
generic = kernel_build_cache._build_metadata(tmp_path)
assert explicit_hopper["cache_key"] == detected_hopper["cache_key"]
assert explicit_hopper["build"]["torch_cuda_arch_list"] == "9.0a"
assert detected_hopper["build"]["torch_cuda_arch_list"] == ""
assert detected_l40s["cache_key"] != detected_hopper["cache_key"]
assert explicit["cache_key"] == detected["cache_key"]
assert explicit["build"]["torch_cuda_arch_list"] == arch
assert detected["build"]["torch_cuda_arch_list"] == ""
assert generic["cache_key"] != detected["cache_key"]
@pytest.mark.parametrize("capability, expected", [((10, 0), "10.0a"), ((10, 3), "10.3a")])
def test_blackwell_detection_reaches_kernel_build(monkeypatch, tmp_path, capability, expected) -> None:
_patch_stable_metadata(monkeypatch)
_patch_cuda_capability(monkeypatch, capability)
metadata = kernel_build_cache._build_metadata(tmp_path)
build_arches = []
def fake_run(command, *, cwd, env):
assert command[:2] == ["./build.sh", "--wheel-dir"]
assert cwd == tmp_path / "fastvideo-kernel"
build_arches.append(env["TORCH_CUDA_ARCH_LIST"])
_write_wheel(Path(command[2]) / WHEEL_NAME)
return ""
monkeypatch.setattr(kernel_build_cache, "_run", fake_run)
wheel = kernel_build_cache._build_wheel(tmp_path, metadata)
try:
assert wheel.is_file()
assert build_arches == [expected]
finally:
shutil.rmtree(wheel.parent)
def test_cache_key_ignores_runtime_only_torch_config(monkeypatch, tmp_path) -> None:
+1
View File
@@ -0,0 +1 @@
"""Trusted, opt-in GPU CI dispatch. Never import this package from a PR checkout."""
+98
View File
@@ -0,0 +1,98 @@
"""Operator CLI. Use the isolated entrypoint from an immutable installation."""
from __future__ import annotations
import argparse
import json
import os
import subprocess
import sys
import threading
from pathlib import Path
from typing import Any
from .dispatcher import (cancellation_signals, execute, make_store, recover,
run_directory, upload_artifacts)
from .pipeline import render_pipeline
from .policy import DEFAULT_REPOSITORY, backend_name, build_request
ROOT = Path(__file__).resolve().parents[2]
def load_config(path: Path) -> dict[str, Any]:
config = json.loads(path.read_text())
if not isinstance(config, dict):
raise ValueError("Operator config must be a JSON object")
for key in ("state_path", "artifacts_dir"):
value = config.get(key)
if not isinstance(value, str) or not Path(value).is_absolute():
raise ValueError(f"{key} must be an absolute local path")
for key, default in (("max_active_prs", 2), ("max_gpus_per_pr", 4), ("max_gpus", 8),
("queue_timeout_seconds", 21600)):
value = config.get(key, default)
if type(value) is not int or value < 1:
raise ValueError(f"{key} must be a positive integer")
config[key] = value
# The initial rollout deliberately cannot raise the user's agreed ceilings.
for key, ceiling in (("max_active_prs", 2), ("max_gpus_per_pr", 4), ("max_gpus", 8)):
if config[key] > ceiling:
raise ValueError(f"{key} exceeds the reviewed ceiling {ceiling}")
if config.get("repository", DEFAULT_REPOSITORY) != DEFAULT_REPOSITORY:
raise ValueError("This worker installation is allowlisted only for hao-ai-lab/FastVideo")
backend_name(config.get("default_backend"))
return config
def main(argv: list[str] | None = None) -> int:
parser = argparse.ArgumentParser(description=__doc__)
subparsers = parser.add_subparsers(dest="command", required=True)
for name in ("run", "render", "upload", "status", "recover"):
sub = subparsers.add_parser(name)
sub.add_argument("--config", type=Path, required=True)
if name in ("render", "upload"):
sub.add_argument("--canonical", type=Path, default=ROOT / ".buildkite/pipeline.yml")
if name == "recover":
sub.add_argument("--build-id", required=True)
args = parser.parse_args(argv)
try:
config = load_config(args.config)
if args.command == "status":
print(json.dumps(make_store(config).snapshot(), indent=2))
return 0
if args.command == "recover":
recover(config, args.build_id)
return 0
backend = backend_name(os.environ.get("CI_GPU_BACKEND"), config.get("default_backend", "slurm"))
if args.command == "upload" and backend == "slurm":
if config.get("default_backend", "slurm") != "slurm":
raise ValueError("Slurm overrides are disabled after promoting another backend: legacy statuses are unqualified")
legacy = config.get("slurm_uploader")
if (not isinstance(legacy, list) or not legacy or
not all(isinstance(x, str) and x for x in legacy) or not Path(legacy[0]).is_absolute()):
raise ValueError("Configure slurm_uploader as the existing trusted uploader argv")
return subprocess.run(legacy, check=False).returncode
if args.command in ("render", "upload"):
import yaml
canonical = yaml.safe_load(args.canonical.read_text())
rendered = render_pipeline(canonical, backend, os.environ.get("TEST_SCOPE", ""),
config.get("queue", "gpu-ci-dispatch"))
payload = json.dumps(rendered, indent=2) + "\n"
if args.command == "render":
print(payload, end="")
return 0
return subprocess.run(["buildkite-agent", "pipeline", "upload", "--no-interpolation"],
input=payload, text=True, check=False).returncode
request = build_request(os.environ, config, ROOT)
cancel = threading.Event()
with cancellation_signals(cancel):
code = execute(request, config, cancel=cancel)
if config.get("upload_artifacts", True):
upload_artifacts(run_directory(config, request["build_id"]))
return code
except (OSError, ValueError, RuntimeError, KeyError, subprocess.SubprocessError) as error:
print(f"GPU CI dispatcher error: {error}", file=sys.stderr)
return 97
if __name__ == "__main__":
raise SystemExit(main())
+307
View File
@@ -0,0 +1,307 @@
# SPDX-License-Identifier: Apache-2.0
"""Persistent GPU admission for cooperating processes on one trusted host.
Keep this SQLite database on local disk, not a shared/network filesystem.
Reservations include backend handles before external creation. Neither stale
heartbeats nor cancellation free active reservations: the caller must stop the
backend worker, confirm it stopped, and then call ``release``. The dispatcher
must also serialize creation and recovery for each build outside this ledger.
"""
from __future__ import annotations
import json
import sqlite3
import time
from collections.abc import Iterator
from contextlib import contextmanager
from pathlib import Path
from typing import Any
class AdmissionStore:
"""Bound distinct PRs and GPU reservations across all selected backends.
``pr_key`` identifies a repository and PR together. ``None`` or an empty
string gives a non-PR build its own slot. Once admitted, all unfinished
builds for a PR share and retain its slot, even between lanes. Completed
build and lane records remain as tombstones; retries must use new lane IDs.
"""
def __init__(self, path: str | Path, max_prs: int = 2, max_gpus_per_pr: int = 4, max_gpus: int = 8):
for name, value in (("max_prs", max_prs), ("max_gpus_per_pr", max_gpus_per_pr), ("max_gpus", max_gpus)):
if type(value) is not int or value < 1:
raise ValueError(f"{name} must be a positive integer")
if str(path) == ":memory:":
raise ValueError("Admission requires a persistent database on local disk")
self.path = Path(path)
self.path.parent.mkdir(parents=True, exist_ok=True)
self.max_prs = max_prs
self.max_gpus_per_pr = max_gpus_per_pr
self.max_gpus = max_gpus
with self._transaction() as connection:
connection.execute("""
CREATE TABLE IF NOT EXISTS settings (
id INTEGER PRIMARY KEY CHECK (id = 1),
schema_version INTEGER NOT NULL,
max_prs INTEGER NOT NULL,
max_gpus_per_pr INTEGER NOT NULL,
max_gpus INTEGER NOT NULL
)
""")
connection.execute("INSERT OR IGNORE INTO settings VALUES (1, 1, ?, ?, ?)",
(max_prs, max_gpus_per_pr, max_gpus))
settings = connection.execute("SELECT * FROM settings WHERE id = 1").fetchone()
expected = (1, max_prs, max_gpus_per_pr, max_gpus)
actual = tuple(settings[name] for name in ("schema_version", "max_prs", "max_gpus_per_pr", "max_gpus"))
if actual != expected:
raise ValueError("Admission database schema or limits conflict with this process")
connection.execute("""
CREATE TABLE IF NOT EXISTS builds (
sequence INTEGER PRIMARY KEY AUTOINCREMENT,
build_id TEXT NOT NULL UNIQUE,
pr_key TEXT,
workload_key TEXT NOT NULL,
backend TEXT NOT NULL,
commit_sha TEXT NOT NULL,
state TEXT NOT NULL CHECK (state IN ('queued', 'admitted', 'finished')),
cancelled INTEGER NOT NULL DEFAULT 0 CHECK (cancelled IN (0, 1)),
created_at REAL NOT NULL,
heartbeat_at REAL NOT NULL,
admitted_at REAL,
finished_at REAL
)
""")
connection.execute("""
CREATE TABLE IF NOT EXISTS allocations (
sequence INTEGER PRIMARY KEY AUTOINCREMENT,
build_id TEXT NOT NULL REFERENCES builds(build_id),
lane_id TEXT NOT NULL,
gpus INTEGER NOT NULL CHECK (gpus > 0),
handle_json TEXT NOT NULL,
state TEXT NOT NULL CHECK (state IN ('waiting', 'active', 'released', 'cancelled')),
created_at REAL NOT NULL,
acquired_at REAL,
released_at REAL,
UNIQUE (build_id, lane_id)
)
""")
connection.execute("CREATE INDEX IF NOT EXISTS builds_workload ON builds(workload_key, state)")
connection.execute("CREATE INDEX IF NOT EXISTS allocations_state ON allocations(state, sequence)")
@contextmanager
def _transaction(self) -> Iterator[sqlite3.Connection]:
connection = sqlite3.connect(str(self.path), timeout=30, isolation_level=None)
connection.row_factory = sqlite3.Row
try:
connection.execute("PRAGMA foreign_keys = ON")
connection.execute("BEGIN IMMEDIATE")
yield connection
connection.commit()
except BaseException:
connection.rollback()
raise
finally:
connection.close()
@staticmethod
def _identifier(value: str, name: str) -> None:
if not isinstance(value, str) or not value.strip():
raise ValueError(f"{name} must be a nonempty string")
@staticmethod
def _build(connection: sqlite3.Connection, build_id: str) -> sqlite3.Row:
build = connection.execute("SELECT * FROM builds WHERE build_id = ?", (build_id, )).fetchone()
if build is None:
raise KeyError(f"Unknown build: {build_id}")
return build
def register(self, build_id: str, pr_key: str | None, backend: str, commit: str) -> None:
"""Register once; an identical duplicate never resets state or queue age."""
for name, value in (("build_id", build_id), ("backend", backend), ("commit", commit)):
self._identifier(value, name)
if pr_key is not None and not isinstance(pr_key, str):
raise ValueError("pr_key must be a string or None")
pr_key = pr_key or None
if pr_key is not None:
self._identifier(pr_key, "pr_key")
workload_key = f"pr:{pr_key}" if pr_key is not None else f"build:{build_id}"
with self._transaction() as connection:
existing = connection.execute("SELECT * FROM builds WHERE build_id = ?", (build_id, )).fetchone()
if existing is not None:
if (existing["pr_key"], existing["backend"], existing["commit_sha"]) != (pr_key, backend, commit):
raise ValueError(f"Build identity changed: {build_id}")
return
active = connection.execute("SELECT 1 FROM builds WHERE workload_key = ? AND state = 'admitted'",
(workload_key, )).fetchone()
now = time.time()
connection.execute("""
INSERT INTO builds (build_id, pr_key, workload_key, backend, commit_sha,
state, created_at, heartbeat_at, admitted_at)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)
""", (build_id, pr_key, workload_key, backend, commit,
"admitted" if active else "queued", now, now, now if active else None))
def try_admit(self, build_id: str) -> bool:
"""Admit the oldest waiting PRs, reserving space for earlier waiters."""
with self._transaction() as connection:
build = self._build(connection, build_id)
if build["cancelled"] or build["state"] == "finished":
return False
if build["state"] == "admitted":
return True
active_count = connection.execute(
"SELECT COUNT(DISTINCT workload_key) FROM builds WHERE state = 'admitted'").fetchone()[0]
available = self.max_prs - active_count
if available <= 0:
return False
eligible = connection.execute("""
SELECT workload_key FROM builds WHERE state = 'queued' AND cancelled = 0
GROUP BY workload_key ORDER BY MIN(sequence) LIMIT ?
""", (available, )).fetchall()
if build["workload_key"] not in {row["workload_key"] for row in eligible}:
return False
connection.execute("""
UPDATE builds SET state = 'admitted', admitted_at = ?
WHERE workload_key = ? AND state = 'queued' AND cancelled = 0
""", (time.time(), build["workload_key"]))
return True
def try_acquire(self, build_id: str, lane_id: str, gpus: int, handle: dict[str, Any]) -> bool:
"""Persist a handle and reserve GPUs before the caller creates a worker.
Pending requests keep FIFO order. A PR whose oldest request cannot fit
its own budget does not block the other PR. Once a request fits its PR
budget it waits for global capacity without smaller requests bypassing
it. Repeating an active request returns True with the same reservation;
released/cancelled lane IDs cannot acquire again.
"""
self._identifier(lane_id, "lane_id")
if type(gpus) is not int or not 1 <= gpus <= min(self.max_gpus_per_pr, self.max_gpus):
raise ValueError("Requested GPUs must fit both the per-PR and total GPU limits")
if not isinstance(handle, dict):
raise ValueError("Backend handle must be a JSON object")
try:
encoded_handle = json.dumps(handle, sort_keys=True, separators=(",", ":"), allow_nan=False)
except (TypeError, ValueError) as error:
raise ValueError("Backend handle must be JSON serializable") from error
with self._transaction() as connection:
build = self._build(connection, build_id)
allocation = connection.execute("SELECT * FROM allocations WHERE build_id = ? AND lane_id = ?",
(build_id, lane_id)).fetchone()
if allocation is not None:
if allocation["gpus"] != gpus or allocation["handle_json"] != encoded_handle:
raise ValueError(f"Lane reservation identity changed: {build_id}/{lane_id}")
if build["cancelled"] or build["state"] != "admitted":
return False
if allocation is not None and allocation["state"] != "waiting":
return allocation["state"] == "active"
if allocation is None:
connection.execute("""
INSERT INTO allocations (build_id, lane_id, gpus, handle_json, state, created_at)
VALUES (?, ?, ?, ?, 'waiting', ?)
""", (build_id, lane_id, gpus, encoded_handle, time.time()))
usage = connection.execute("""
SELECT b.workload_key, SUM(a.gpus) AS gpus FROM allocations a
JOIN builds b ON b.build_id = a.build_id WHERE a.state = 'active'
GROUP BY b.workload_key
""").fetchall()
gpu_usage = {row["workload_key"]: row["gpus"] for row in usage}
total_gpus = sum(gpu_usage.values())
waiting = connection.execute("""
SELECT a.*, b.workload_key FROM allocations a
JOIN builds b ON b.build_id = a.build_id
WHERE a.state = 'waiting' AND b.state = 'admitted' AND b.cancelled = 0
ORDER BY a.sequence
""").fetchall()
seen_workloads = set()
for request in waiting:
workload_key = request["workload_key"]
if workload_key in seen_workloads:
continue
seen_workloads.add(workload_key)
if gpu_usage.get(workload_key, 0) + request["gpus"] > self.max_gpus_per_pr:
continue
if (request["build_id"], request["lane_id"]) != (build_id, lane_id):
return False
if total_gpus + gpus > self.max_gpus:
return False
connection.execute("""
UPDATE allocations SET state = 'active', acquired_at = ?
WHERE build_id = ? AND lane_id = ?
""", (time.time(), build_id, lane_id))
return True
return False
def release(self, build_id: str, lane_id: str) -> None:
"""Release only after the caller confirms the worker has stopped.
Missing/already released reservations are harmless for reconciliation.
This method does not stop workers or infer their status from timestamps.
"""
with self._transaction() as connection:
self._build(connection, build_id)
connection.execute("""
UPDATE allocations SET state = 'released', released_at = ?
WHERE build_id = ? AND lane_id = ? AND state IN ('active', 'waiting')
""", (time.time(), build_id, lane_id))
def finish(self, build_id: str) -> None:
"""Finish or cancel a queued build, refusing any active reservations."""
with self._transaction() as connection:
build = self._build(connection, build_id)
if build["state"] == "finished":
return
active = connection.execute("SELECT 1 FROM allocations WHERE build_id = ? AND state = 'active'",
(build_id, )).fetchone()
if active is not None:
raise RuntimeError(f"Cannot finish build with active GPU allocations: {build_id}")
now = time.time()
connection.execute("""
UPDATE allocations SET state = 'cancelled', released_at = ?
WHERE build_id = ? AND state = 'waiting'
""", (now, build_id))
connection.execute("UPDATE builds SET state = 'finished', finished_at = ? WHERE build_id = ?",
(now, build_id))
def heartbeat(self, build_id: str) -> None:
"""Record liveness for operators; no timeout ever frees a reservation."""
with self._transaction() as connection:
self._build(connection, build_id)
connection.execute("UPDATE builds SET heartbeat_at = ? WHERE build_id = ? AND state != 'finished'",
(time.time(), build_id))
def request_cancel(self, build_id: str) -> None:
"""Block further admission/launch reservations without freeing GPUs."""
with self._transaction() as connection:
self._build(connection, build_id)
connection.execute("UPDATE builds SET cancelled = 1 WHERE build_id = ?", (build_id, ))
connection.execute("""
UPDATE allocations SET state = 'cancelled', released_at = ?
WHERE build_id = ? AND state = 'waiting'
""", (time.time(), build_id))
def is_cancelled(self, build_id: str) -> bool:
with self._transaction() as connection:
return bool(self._build(connection, build_id)["cancelled"])
def snapshot(self) -> dict[str, Any]:
"""Return a consistent reconciliation view, including terminal records."""
with self._transaction() as connection:
builds = [dict(row) for row in connection.execute("SELECT * FROM builds ORDER BY sequence")]
allocations = [dict(row) for row in connection.execute("SELECT * FROM allocations ORDER BY sequence")]
for build in builds:
build["commit"] = build.pop("commit_sha")
build["cancelled"] = bool(build["cancelled"])
for allocation in allocations:
allocation["handle"] = json.loads(allocation.pop("handle_json"))
return {
"limits": {
"max_prs": self.max_prs,
"max_gpus_per_pr": self.max_gpus_per_pr,
"max_gpus": self.max_gpus,
},
"builds": builds,
"allocations": allocations,
}
+395
View File
@@ -0,0 +1,395 @@
"""Trusted execution adapters. No PR code executes in the controller process.
Modal API: https://modal.com/docs/sdk/py/latest/Sandbox
The controller persists deterministic handles before calling ``run``.
"""
from __future__ import annotations
import hashlib
import importlib
import json
import re
import subprocess
import threading
import time
from pathlib import Path, PurePosixPath
from typing import Any
REPOSITORY = "https://github.com/hao-ai-lab/FastVideo.git"
ENV_ALLOWLIST = {
"FASTVIDEO_GOLDEN_TEST_FILES", "FASTVIDEO_SSIM_TEST_FILES", "BUILDKITE_BRANCH",
"BUILDKITE_SOURCE", "BUILDKITE_PULL_REQUEST", "BUILDKITE_BUILD_NUMBER", "BUILDKITE_BUILD_URL",
"BUILDKITE_COMMIT", "BUILDKITE_REPO", "BUILDKITE_BUILD_ID", "BUILDKITE_JOB_ID",
}
def _image(value: str) -> str:
if not re.fullmatch(r"[A-Za-z0-9._/:\-]+@sha256:[0-9a-f]{64}", value):
raise ValueError("Backend image must be pinned by sha256 digest")
return value
def _subpath(value: str) -> str:
path = PurePosixPath(value)
if not value or path.is_absolute() or ".." in path.parts or value in {".", "/"}:
raise ValueError("PVC mounts require a dedicated relative CI subpath")
return value
class Backend:
name = "abstract"
def __init__(self, config: dict, trusted_dir: Path, run_dir: Path):
self.config = config
self.trusted_dir = trusted_dir
self.run_dir = run_dir
self.image = _image(config["image"])
self.run_dir.mkdir(parents=True, exist_ok=True)
def handle(self, build_id: str, lane_id: str) -> dict:
identity = hashlib.sha256(f"{build_id}\0{lane_id}".encode()).hexdigest()[:32]
return {"backend": self.name, "name": f"fv-ci-{identity}", "run_id": identity}
def _environment(self, request: dict, lane: dict) -> dict[str, str]:
if request["repository"] != REPOSITORY or not re.fullmatch(r"[0-9a-f]{40}", request["commit"]):
raise ValueError("Worker requires the fixed repository and an immutable commit")
if type(lane["gpus"]) is not int or lane["gpus"] not in range(1, 5):
raise ValueError("A lane must request between one and four GPUs")
script = lane["script"]
if script != ".buildkite/scripts/unit_test.sh" and not re.fullmatch(
r"\.buildkite/scripts/lanes/[a-z_]+\.sh", script):
raise ValueError("Invalid lane script")
extras = lane.get("extras", ["test"])
if not extras or any(not re.fullmatch(r"[a-z][a-z0-9-]*", extra) for extra in extras):
raise ValueError("Invalid dependency extras")
if str(lane["fa4"]) not in {"0", "1"}:
raise ValueError("Invalid attention policy")
build_id = str(request.get("buildkite_build_id") or request["build_id"])
job_id = str(request.get("buildkite_job_id") or request.get("env", {}).get("BUILDKITE_JOB_ID")
or request["build_id"])
if any(not re.fullmatch(r"[A-Za-z0-9][A-Za-z0-9._-]{0,199}", value) for value in (build_id, job_id)):
raise ValueError("Invalid Buildkite build/job identity")
pr_number = str(request.get("pr_number") or "false")
if pr_number != "false" and not re.fullmatch(r"[1-9][0-9]*", pr_number):
raise ValueError("Invalid PR number")
env = {key: str(value) for key, value in request.get("env", {}).items() if key in ENV_ALLOWLIST}
env.update({
"BUILDKITE_COMMIT": request["commit"], "BUILDKITE_REPO": REPOSITORY,
"BUILDKITE_BUILD_ID": build_id, "BUILDKITE_JOB_ID": job_id,
"BUILDKITE_PULL_REQUEST": pr_number, "FASTVIDEO_CONTAINER_IMAGE_REF": self.image,
"FASTVIDEO_CI_REPOSITORY": REPOSITORY, "FASTVIDEO_CI_COMMIT": request["commit"],
"FASTVIDEO_CI_PR_NUMBER": str(request.get("pr_number") or ""),
"FASTVIDEO_CI_GPUS": str(lane["gpus"]), "FASTVIDEO_CI_SCRIPT": script,
"FASTVIDEO_CI_EXTRAS": ",".join(extras), "FASTVIDEO_FA4": str(lane["fa4"]),
"FASTVIDEO_CI_KERNEL": "1" if lane.get("kernel", True) else "0",
"FASTVIDEO_CI_LOCAL_ONLY": "1", "FASTVIDEO_SSIM_BOOTSTRAP_MODE": "0",
"TEST_SCOPE": request.get("scope", "direct"),
})
return env
def run(self, handle: dict, request: dict, lane: dict, cancel: threading.Event) -> int:
raise NotImplementedError
def stop(self, handle: dict) -> bool:
raise NotImplementedError
class KubernetesBackend(Backend):
name = "kubernetes"
def __init__(self, config: dict, trusted_dir: Path, run_dir: Path):
super().__init__(config, trusted_dir, run_dir)
if config.get("namespace", "vllm") != "vllm":
raise ValueError("Kubernetes CI is restricted to namespace vllm")
for kind in ("cache", "artifacts"):
claim = config.get(f"{kind}_pvc")
if claim:
if claim in {"lustre-pvc-vllm", "nfs-pvc-vllm"}:
raise ValueError("Personal development PVCs cannot be mounted in CI")
_subpath(config.get(f"{kind}_subpath", ""))
def handle(self, build_id: str, lane_id: str) -> dict:
result = super().handle(build_id, lane_id)
result.update({"namespace": "vllm", "context": self.config.get("context")})
return result
def _command(self, handle: dict, *args: str) -> list[str]:
if handle.get("namespace") != "vllm" or handle.get("context") != self.config.get("context"):
raise ValueError("Handle does not match configured Kubernetes target")
command = ["kubectl", "--request-timeout=20s", "--namespace=vllm"]
if handle.get("context"):
command += ["--context", handle["context"]]
return [*command, *args]
def _kubectl(self, handle: dict, *args: str, body: dict | None = None) -> str:
result = subprocess.run(self._command(handle, *args), input=json.dumps(body) if body else None,
text=True, capture_output=True, timeout=35, check=False)
if result.returncode:
raise RuntimeError(f"kubectl {args[0]} failed: {result.stderr.strip()}")
return result.stdout
def _job(self, handle: dict) -> dict | None:
raw = self._kubectl(handle, "get", "job", handle["name"], "--ignore-not-found", "-o", "json")
return json.loads(raw) if raw.strip() else None
def _pods(self, handle: dict) -> list[dict]:
raw = self._kubectl(handle, "get", "pods", "-l", f"job-name={handle['name']}", "-o", "json")
return json.loads(raw)["items"]
def manifest(self, handle: dict, request: dict, lane: dict) -> dict:
env = [{"name": key, "value": value} for key, value in self._environment(request, lane).items()]
if self.config.get("hf_secret"):
env.append({"name": "HF_TOKEN", "valueFrom": {"secretKeyRef": {
"name": self.config["hf_secret"], "key": self.config.get("hf_secret_key", "HF_TOKEN")}}})
volumes = [{"name": "workspace", "emptyDir": {}}, {"name": "shm", "emptyDir": {
"medium": "Memory", "sizeLimit": self.config.get("shared_memory", "8Gi")}}]
mounts = [{"name": "workspace", "mountPath": "/workspace"}, {"name": "shm", "mountPath": "/dev/shm"}]
for kind, path in (("cache", "/ci-cache"), ("artifacts", "/workspace/artifacts")):
if self.config.get(f"{kind}_pvc"):
subpath = _subpath(self.config[f"{kind}_subpath"])
if kind == "artifacts":
subpath += "/" + handle["name"]
volumes.append({"name": kind, "persistentVolumeClaim": {"claimName": self.config[f"{kind}_pvc"]}})
mounts.append({"name": kind, "mountPath": path, "subPath": subpath, "readOnly": kind == "cache"})
resources = {"cpu": str(self.config.get("cpu", "8")), "memory": self.config.get("memory", "64Gi"),
"nvidia.com/gpu": str(lane["gpus"])}
labels = {"app.kubernetes.io/managed-by": "fastvideo-gpu-ci", "fastvideo-ci/run-id": handle["run_id"]}
manifest = {"apiVersion": "batch/v1", "kind": "Job", "metadata": {
"name": handle["name"], "namespace": "vllm", "labels": labels}, "spec": {
"backoffLimit": 0, "activeDeadlineSeconds": int(lane["timeout_seconds"]),
"ttlSecondsAfterFinished": 86400, "template": {"metadata": {"labels": labels}, "spec": {
"restartPolicy": "Never", "automountServiceAccountToken": False,
"hostNetwork": False, "hostPID": False, "hostIPC": False,
"nodeSelector": {"kubernetes.io/arch": "arm64", "nvidia.com/gpu.product": "NVIDIA-GB200"},
"tolerations": [{"key": "nvidia.com/gpu", "operator": "Exists", "effect": "NoSchedule"}],
"containers": [{"name": "worker", "image": self.image,
"command": ["/bin/bash", "-c", (self.trusted_dir / "worker.sh").read_text()],
"env": env, "resources": {"requests": resources, "limits": resources},
"securityContext": {"privileged": False, "allowPrivilegeEscalation": False,
"capabilities": {"drop": ["ALL"], "add": [
"CHOWN", "DAC_OVERRIDE", "FOWNER", "SETGID", "SETUID"]}},
"volumeMounts": mounts}], "volumes": volumes}}}}
if self.config.get("artifacts_pvc"):
# Only this fixed init command sees the CI artifact parent. PR code
# mounts its own run directory and cannot alter adjacent runs.
manifest["spec"]["template"]["spec"]["initContainers"] = [{
"name": "prepare-artifacts", "image": self.image,
"command": ["mkdir", "-p", "/ci-artifacts/" + handle["name"]],
"securityContext": {"allowPrivilegeEscalation": False, "capabilities": {"drop": ["ALL"]}},
"volumeMounts": [{"name": "artifacts", "mountPath": "/ci-artifacts",
"subPath": _subpath(self.config["artifacts_subpath"])}]}]
return manifest
@staticmethod
def _owned(job: dict, handle: dict) -> bool:
return job.get("metadata", {}).get("labels", {}).get("fastvideo-ci/run-id") == handle["run_id"]
def run(self, handle: dict, request: dict, lane: dict, cancel: threading.Event) -> int:
if self._job(handle) is not None or self._pods(handle):
raise RuntimeError("Resource already exists; explicit reconciliation is required")
manifest = self.manifest(handle, request, lane)
(self.run_dir / f"{handle['name']}.job.json").write_text(json.dumps(manifest, indent=2))
if not self.config.get("artifacts_pvc"):
(self.run_dir / f"{handle['name']}.artifacts.txt").write_text(
"No dedicated artifacts PVC configured. Controller logs/status survive; worker artifacts are ephemeral.\n")
if cancel.is_set():
return 130
marker = self.run_dir / f"{handle['name']}.creating"
marker.touch()
self._kubectl(handle, "create", "-f", "-", body=manifest)
marker.unlink()
deadline = time.monotonic() + int(lane["timeout_seconds"])
stream = None
log_file = (self.run_dir / f"{handle['name']}.log").open("w")
try:
while True:
if cancel.is_set() or time.monotonic() >= deadline:
if not self.stop(handle):
raise RuntimeError("Cannot confirm Kubernetes job cancellation")
return 130 if cancel.is_set() else 124
job = self._job(handle)
if job is None or not self._owned(job, handle):
raise RuntimeError("Kubernetes job disappeared or its ownership changed")
pods = self._pods(handle)
if len(pods) > 1:
raise RuntimeError("Unexpected replacement or duplicate CI pods")
for pod in pods:
owners = pod.get("metadata", {}).get("ownerReferences", [])
if not any(owner.get("uid") == job["metadata"]["uid"] for owner in owners):
raise RuntimeError("Pod does not belong to the expected job")
statuses = pod.get("status", {}).get("containerStatuses", [])
worker = next((status for status in statuses if status["name"] == "worker"), {})
state = worker.get("state", {})
if stream is None and ("running" in state or "terminated" in state):
stream = subprocess.Popen(self._command(handle, "logs", "-f", pod["metadata"]["name"],
"-c", "worker", "--timestamps"),
stdout=log_file, stderr=subprocess.STDOUT, text=True)
terminated = state.get("terminated")
if terminated is not None:
exit_code = terminated.get("exitCode")
if type(exit_code) is not int:
raise RuntimeError("Worker termination has no numeric exit status")
(self.run_dir / f"{handle['name']}.status.json").write_text(json.dumps(pod["status"], indent=2))
# Streaming may fail independently of the pod. Fetch a
# final complete log before accepting even a zero exit.
final_log = self._kubectl(handle, "logs", pod["metadata"]["name"], "-c", "worker", "--timestamps")
if not final_log.strip():
raise RuntimeError("Worker completed but Kubernetes logs are empty")
(self.run_dir / f"{handle['name']}.final.log").write_text(final_log)
return exit_code
if pod.get("status", {}).get("phase") == "Failed":
raise RuntimeError("Pod failed without a worker exit status")
if any(condition.get("type") == "Failed" and condition.get("status") == "True"
for condition in job.get("status", {}).get("conditions", [])):
raise RuntimeError("Kubernetes job failed before worker completion")
cancel.wait(2)
finally:
if stream is not None:
try:
stream.wait(timeout=10)
except subprocess.TimeoutExpired:
stream.terminate()
try:
stream.wait(timeout=5)
except subprocess.TimeoutExpired:
stream.kill()
stream.wait()
log_file.close()
def stop(self, handle: dict) -> bool:
job = self._job(handle)
if job is not None:
if not self._owned(job, handle):
return False
self._kubectl(handle, "delete", "job", handle["name"], "--cascade=foreground",
"--wait=true", "--timeout=20s")
(self.run_dir / f"{handle['name']}.creating").unlink(missing_ok=True)
elif (self.run_dir / f"{handle['name']}.creating").exists():
# A timed-out create may still commit remotely. Keep its GPU lease.
return False
# A missing Job alone is insufficient: orphaned Pods may still use GPUs.
return self._job(handle) is None and not self._pods(handle)
class ModalBackend(Backend):
name = "modal"
def __init__(self, config: dict, trusted_dir: Path, run_dir: Path):
super().__init__(config, trusted_dir, run_dir)
if config.get("gpu", "L40S") != "L40S":
raise ValueError("Modal default must remain L40S; Hopper lanes are selected by reviewed lane policy")
self._sandboxes: dict[str, Any] = {}
@staticmethod
def _modal() -> Any:
return importlib.import_module("modal")
@staticmethod
def _gpu_type(lane: dict) -> str:
gpu_type = lane.get("modal_gpu", "L40S")
if gpu_type not in ("L40S", "H100"):
raise ValueError("Reviewed Modal hardware must be L40S or H100")
return gpu_type
def _environment(self, request: dict, lane: dict) -> dict[str, str]:
self._gpu_type(lane)
# Attention/reference policy differs from the GB200 lanes. Copy the
# lane record so concurrent tasks never mutate shared policy or image.
modal_lane = {**lane, "fa4": lane.get("modal_fa4", lane["fa4"])}
return super()._environment(request, modal_lane)
def handle(self, build_id: str, lane_id: str) -> dict:
result = super().handle(build_id, lane_id)
result["app"] = self.config.get("app", "fastvideo-gpu-ci")
return result
def _lookup(self, handle: dict) -> Any:
if handle["app"] != self.config.get("app", "fastvideo-gpu-ci"):
raise ValueError("Handle does not match configured Modal app")
modal = self._modal()
try:
return modal.Sandbox.from_name(handle["app"], handle["name"])
except modal.exception.NotFoundError:
return None
def run(self, handle: dict, request: dict, lane: dict, cancel: threading.Event) -> int:
env = self._environment(request, lane)
if self._lookup(handle) is not None:
raise RuntimeError("Sandbox already exists; explicit reconciliation is required")
modal = self._modal()
image = modal.Image.from_registry(self.image).add_local_file(
str(self.trusted_dir / "worker.sh"), "/opt/fastvideo-ci/worker.sh")
app = modal.App.lookup(handle["app"], create_if_missing=True)
secrets = [modal.Secret.from_name(self.config["hf_secret"])] if self.config.get("hf_secret") else []
(self.run_dir / f"{handle['name']}.artifacts.txt").write_text(
"Modal adapter retains controller logs/status. Worker artifacts are ephemeral.\n")
if cancel.is_set():
return 130
marker = self.run_dir / f"{handle['name']}.creating"
marker.touch()
sandbox = modal.Sandbox.create(
"/bin/bash", "/opt/fastvideo-ci/worker.sh", app=app, name=handle["name"], image=image,
env=env, secrets=secrets, gpu=f"{self._gpu_type(lane)}:{lane['gpus']}",
timeout=int(lane["timeout_seconds"]),
cpu=float(self.config.get("cpu", 8)), memory=int(self.config.get("memory_mib", 65536)),
tags={"fastvideo-ci-run-id": handle["run_id"]})
self._sandboxes[handle["name"]] = sandbox
marker.unlink()
(self.run_dir / f"{handle['name']}.sandbox.json").write_text(json.dumps({"id": sandbox.object_id}))
return self._watch(sandbox, handle, lane, cancel)
def _watch(self, sandbox: Any, handle: dict, lane: dict, cancel: threading.Event) -> int:
errors: list[Exception] = []
def capture(stream: Any, suffix: str) -> None:
try:
with (self.run_dir / f"{handle['name']}.{suffix}.log").open("w") as output:
for line in stream:
output.write(line)
output.flush()
except Exception as error:
errors.append(error)
threads = [threading.Thread(target=capture, args=(sandbox.stdout, "stdout"), daemon=True),
threading.Thread(target=capture, args=(sandbox.stderr, "stderr"), daemon=True)]
for thread in threads:
thread.start()
deadline = time.monotonic() + int(lane["timeout_seconds"])
try:
while True:
if cancel.is_set() or time.monotonic() >= deadline:
if not self.stop(handle):
raise RuntimeError("Cannot confirm Modal sandbox cancellation")
return 130 if cancel.is_set() else 124
code = sandbox.poll()
if code is not None:
if type(code) is not int:
raise RuntimeError("Sandbox returned a nonnumeric exit status")
(self.run_dir / f"{handle['name']}.status.json").write_text(json.dumps({"exit_code": code}))
for thread in threads:
thread.join(timeout=5)
if errors or any(thread.is_alive() for thread in threads):
raise RuntimeError("Modal log collection did not complete")
return code
cancel.wait(2)
finally:
# Do not detach a running sandbox: root must confirm stop before
# freeing its shared GPU lease, including after SDK/network errors.
for thread in threads:
thread.join(timeout=1)
def stop(self, handle: dict) -> bool:
sandbox = self._sandboxes.get(handle["name"]) or self._lookup(handle)
marker = self.run_dir / f"{handle['name']}.creating"
if sandbox is None:
return not marker.exists()
if sandbox.get_tags().get("fastvideo-ci-run-id") != handle["run_id"]:
return False
code = sandbox.terminate(wait=True)
confirmed = type(code) is int and sandbox.poll() is not None
if confirmed:
marker.unlink(missing_ok=True)
sandbox.detach()
self._sandboxes.pop(handle["name"], None)
return confirmed
+25
View File
@@ -0,0 +1,25 @@
{
"repository": "https://github.com/hao-ai-lab/FastVideo.git",
"default_backend": "slurm",
"slurm_uploader": ["/opt/fastvideo-ci-runner/pipeline-upload"],
"queue": "gpu-ci-dispatch",
"state_path": "/var/lib/fastvideo-gpu-ci/admission.sqlite3",
"artifacts_dir": "/var/lib/fastvideo-gpu-ci/runs",
"max_active_prs": 2,
"max_gpus_per_pr": 4,
"max_gpus": 8,
"queue_timeout_seconds": 21600,
"upload_artifacts": true,
"kubernetes": {
"namespace": "vllm",
"image": "ghcr.io/hao-ai-lab/fastvideo/fastvideo-dev@sha256:REPLACE_WITH_REVIEWED_ARM64_SM100_DIGEST",
"cpu": "8",
"memory": "64Gi",
"shared_memory": "8Gi"
},
"modal": {
"app": "fastvideo-gpu-ci",
"gpu": "L40S",
"image": "ghcr.io/hao-ai-lab/fastvideo/fastvideo-dev@sha256:REPLACE_WITH_REVIEWED_AMD64_SM89_DIGEST"
}
}
+217
View File
@@ -0,0 +1,217 @@
"""Suite lifecycle and local-host recovery for opt-in Buildkite GPU backends."""
from __future__ import annotations
import fcntl
import hashlib
import json
import signal
import subprocess
import threading
import time
from concurrent.futures import Future, ThreadPoolExecutor
from contextlib import contextmanager
from pathlib import Path
from typing import Any, Iterator
from .admission import AdmissionStore
INFRASTRUCTURE_FAILURE = 97
def run_directory(config: dict[str, Any], build_id: str) -> Path:
digest = hashlib.sha256(build_id.encode()).hexdigest()[:32]
return Path(config["artifacts_dir"]) / digest
@contextmanager
def build_lock(directory: Path) -> Iterator[None]:
"""Serialize launch and recovery; never stop a resource during its creation."""
directory.mkdir(parents=True, exist_ok=True, mode=0o700)
with (directory / "dispatcher.lock").open("a") as lock:
try:
fcntl.flock(lock.fileno(), fcntl.LOCK_EX | fcntl.LOCK_NB)
except BlockingIOError as error:
raise RuntimeError("A dispatcher is still running; cancel its Buildkite job first") from error
try:
yield
finally:
fcntl.flock(lock.fileno(), fcntl.LOCK_UN)
def make_store(config: dict[str, Any]) -> AdmissionStore:
return AdmissionStore(config["state_path"], config.get("max_active_prs", 2),
config.get("max_gpus_per_pr", 4), config.get("max_gpus", 8))
def make_backend(config: dict[str, Any], backend: str, directory: Path) -> Any:
from .backends import KubernetesBackend, ModalBackend
cls = {"vllm": KubernetesBackend, "modal": ModalBackend}.get(backend)
if cls is None:
raise ValueError("The new dispatcher supports only modal and vllm")
backend_config = config["kubernetes" if backend == "vllm" else "modal"]
return cls(backend_config, Path(__file__).resolve().parent, directory)
def _run_lane(backend: Any, handle: dict[str, Any], request: dict[str, Any], lane: dict[str, Any],
cancel: threading.Event) -> tuple[int, bool, str]:
error = ""
try:
code = backend.run(handle, request, lane, cancel)
if type(code) is not int or not 0 <= code <= 255:
raise RuntimeError("Backend did not return a valid process exit status")
except Exception as exception:
code, error = INFRASTRUCTURE_FAILURE, str(exception)
try:
stopped = backend.stop(handle)
except Exception as exception:
stopped = False
error += f"; stop failed: {exception}"
if not stopped:
code = INFRASTRUCTURE_FAILURE
error += "; worker termination unconfirmed; reservation retained"
return code, stopped, error
def execute(request: dict[str, Any], config: dict[str, Any], *, backend: Any = None,
cancel: threading.Event | None = None, poll_seconds: float = 2.0) -> int:
"""Run one suite; independent processes share the PR/GPU admission ledger."""
cancel = cancel or threading.Event()
build_id = request["build_id"]
directory = run_directory(config, build_id)
store = make_store(config)
if any(lane["gpus"] > min(store.max_gpus_per_pr, store.max_gpus) for lane in request["lanes"]):
raise ValueError("Selected lane does not fit the configured GPU budget")
backend = backend or make_backend(config, request["backend"], directory)
results: dict[str, Any] = {}
code = 0
with build_lock(directory):
# Never reattach to an ambiguous prior launch. Buildkite retries have
# a new job ID; abandoned attempts need explicit backend reconciliation.
if any(item["build_id"] == build_id for item in store.snapshot()["builds"]):
raise RuntimeError("Build attempt already exists; inspect/recover it instead of replaying it")
store.register(build_id, request["pr_key"], request["backend"], request["commit"])
(directory / "request.json").write_text(json.dumps(request, indent=2) + "\n")
deadline = time.monotonic() + config.get("queue_timeout_seconds", 21600)
# A docs-only merge plan requires no GPUs and must not wait behind
# unrelated PRs. Arrival order is not evidence of commit recency, so
# never cancel other SHAs merely because this request arrived later.
admitted = not request["lanes"]
while not admitted:
store.heartbeat(build_id)
if cancel.is_set() or store.is_cancelled(build_id) or time.monotonic() >= deadline:
cancel.set()
store.request_cancel(build_id)
store.finish(build_id)
code = INFRASTRUCTURE_FAILURE
break
admitted = store.try_admit(build_id)
if not admitted:
cancel.wait(poll_seconds)
pending = list(request["lanes"]) if admitted else []
golden_selected = any(lane["key"] == "golden-gate" for lane in pending)
running: dict[Future[tuple[int, bool, str]], dict[str, Any]] = {}
# Lane concurrency is additionally bounded by transactional GPU leases.
with ThreadPoolExecutor(max_workers=4) as pool:
while pending or running:
store.heartbeat(build_id)
if store.is_cancelled(build_id):
cancel.set()
if cancel.is_set():
store.request_cancel(build_id)
for lane in pending:
results[lane["key"]] = {"state": "cancelled"}
pending.clear()
code = code or INFRASTRUCTURE_FAILURE
for future, lane in list(running.items()):
if not future.done():
continue
lane_code, stopped, error = future.result()
if stopped:
store.release(build_id, lane["key"])
else:
cancel.set()
results[lane["key"]] = {"state": "passed" if lane_code == 0 else "failed",
"exit_code": lane_code, "stopped": stopped, "error": error}
print(f"--- {lane['key']}: exit {lane_code}", flush=True)
code = code or lane_code
del running[future]
for lane in list(pending):
if cancel.is_set() or len(running) >= 4:
break
if golden_selected and not lane["fastcheck"] and lane["key"] != "golden-gate":
golden = results.get("golden-gate")
if golden is None:
continue
if golden["state"] != "passed":
results[lane["key"]] = {"state": "skipped", "reason": "golden gate failed"}
pending.remove(lane)
continue
handle = backend.handle(build_id, lane["key"])
if store.try_acquire(build_id, lane["key"], lane["gpus"], handle):
print(f"--- Starting {lane['key']} on {request['backend']} ({lane['gpus']} GPUs)", flush=True)
future = pool.submit(_run_lane, backend, handle, request, lane, cancel)
running[future] = lane
pending.remove(lane)
if pending or running:
# wait() on an already-set Event spins; cancellation still
# polls bounded backend shutdowns at the normal interval.
time.sleep(poll_seconds)
active = [a for a in store.snapshot()["allocations"] if a["build_id"] == build_id and a["state"] == "active"]
if not active:
store.finish(build_id)
else:
code = INFRASTRUCTURE_FAILURE
summary = {"build_id": build_id, "backend": request["backend"], "commit": request["commit"],
"exit_code": code, "recovery_required": bool(active), "lanes": results}
(directory / "summary.json").write_text(json.dumps(summary, indent=2) + "\n")
print(json.dumps(summary, indent=2), flush=True)
return code
def recover(config: dict[str, Any], build_id: str, *, backend: Any = None) -> None:
"""Stop abandoned workers first, then release reservations. Never TTL-unlock."""
store = make_store(config)
build = next((b for b in store.snapshot()["builds"] if b["build_id"] == build_id), None)
if build is None:
raise ValueError("Unknown build ID")
directory = run_directory(config, build_id)
with build_lock(directory):
store.request_cancel(build_id)
backend = backend or make_backend(config, build["backend"], directory)
uncertain = []
for allocation in store.snapshot()["allocations"]:
if allocation["build_id"] != build_id or allocation["state"] != "active":
continue
try:
stopped = backend.stop(allocation["handle"])
except Exception:
stopped = False
if stopped:
store.release(build_id, allocation["lane_id"])
else:
uncertain.append(allocation["lane_id"])
if uncertain:
raise RuntimeError(f"Could not confirm termination; reservations retained: {uncertain}")
store.finish(build_id)
def upload_artifacts(directory: Path) -> None:
"""Only upload controller-owned logs/metadata, never arbitrary PR globs."""
for path in sorted(directory.iterdir()):
if path.is_file() and not path.is_symlink() and path.suffix in (".json", ".log"):
subprocess.run(["buildkite-agent", "artifact", "upload", path.name], cwd=directory, check=True)
@contextmanager
def cancellation_signals(event: threading.Event) -> Iterator[None]:
previous: dict[signal.Signals, Any] = {}
def handler(signum: int, frame: Any) -> None:
event.set()
for signum in (signal.SIGINT, signal.SIGTERM):
previous[signum] = signal.signal(signum, handler)
try:
yield
finally:
for signum, old_handler in previous.items():
signal.signal(signum, old_handler)
+11
View File
@@ -0,0 +1,11 @@
"""Launch with python -I to exclude job-controlled PYTHONPATH and user packages."""
import sys
from pathlib import Path
sys.path.insert(0, str(Path(__file__).resolve().parents[2]))
from scripts.gpu_ci.__main__ import main # noqa: E402
if __name__ == "__main__":
raise SystemExit(main())
+69
View File
@@ -0,0 +1,69 @@
# SPDX-License-Identifier: Apache-2.0
"""Render an opt-in GPU backend without changing the canonical Slurm graph.
The caller is the operator-owned pipeline uploader. It must load this module
from its reviewed installation, never from a pull request checkout.
"""
from __future__ import annotations
import copy
import re
from typing import Any
BACKENDS = frozenset({"slurm", "modal", "vllm"})
STATUS_SUFFIXES = {
"fastcheck": "fastcheck-passed",
"full": "full-suite-passed",
"merge": "full-suite-passed",
"direct": "direct-test-completed",
"scheduled": "scheduled-ssim-passed",
}
def render_pipeline(
canonical: dict[str, Any],
backend: str,
scope: str,
queue: str = "gpu-ci-dispatch",
) -> dict[str, Any]:
"""Return the original graph or one trusted, backend-specific suite step.
Only enumerated backend and scope values reach the generated step. Request
SHA, PR identity, direct lane, and merge plan remain build metadata and are
validated independently by the trusted dispatcher. Arbitrary canonical
graph commands, environment, hooks, and notifications are not copied into
the new control plane.
"""
backend = backend or "slurm"
scope = scope or "fastcheck"
if backend not in BACKENDS:
raise ValueError(f"Unsupported GPU CI backend: {backend!r}")
if scope not in STATUS_SUFFIXES:
raise ValueError(f"Unsupported GPU CI scope: {scope!r}")
if backend == "slurm":
return copy.deepcopy(canonical)
if not re.fullmatch(r"[a-zA-Z0-9][a-zA-Z0-9_-]{0,63}", queue):
raise ValueError("GPU CI queue must be an operator-configured queue name")
return {
"notify": [{
"github_commit_status": {
"context": f"gpu-ci/{backend}/{STATUS_SUFFIXES[scope]}"
}
}],
"steps": [{
"label": f"GPU CI ({backend}, {scope})",
"key": f"gpu-ci-{backend}-{scope}",
"command": "/opt/fastvideo-gpu-ci/run",
# Admission waiting happens inside this command and uses this
# budget. An unscheduled Buildkite job has no command timeout yet.
"timeout_in_minutes": 480,
"env": {
"CI_GPU_BACKEND": backend,
"TEST_SCOPE": scope,
},
"agents": {
"queue": queue
},
}],
}
+145
View File
@@ -0,0 +1,145 @@
"""Reviewed lane policy and validation at the trusted dispatcher boundary."""
from __future__ import annotations
import re
from pathlib import Path
from typing import Any, Mapping
BACKENDS = ("slurm", "modal", "vllm")
SCOPES = ("fastcheck", "full", "merge", "direct", "scheduled")
DEFAULT_REPOSITORY = "https://github.com/hao-ai-lab/FastVideo.git"
def _lane(key: str, public_type: str, script: str, gpus: int = 1, *, fastcheck: bool = False,
extras: tuple[str, ...] = ("test",), fa4: str = "0") -> dict[str, Any]:
return {"key": key, "public_type": public_type, "script": script, "gpus": gpus,
"fastcheck": fastcheck, "extras": list(extras), "fa4": fa4, "timeout_seconds": 5400,
"kernel": key != "dreamverse",
"modal_gpu": "H100" if key in {"encoder", "kernel-tests", "training-vsa"} else "L40S",
"modal_fa4": "0" if key in {"transformer", "training", "distillation", "self-forcing",
"lora-training", "training-vsa", "train-framework"} else "1"}
LANES = (
_lane("encoder", "encoder", "lanes/encoder.sh", fastcheck=True),
_lane("vae", "vae", "lanes/vae.sh", fastcheck=True),
_lane("transformer", "transformer", "lanes/transformer.sh", fastcheck=True),
_lane("kernel-tests", "kernel_tests", "lanes/kernel_tests.sh", fastcheck=True),
_lane("unit", "unit_test", "unit_test.sh", fastcheck=True),
_lane("dreamverse", "dreamverse_app", "lanes/dreamverse.sh", fastcheck=True,
extras=("test", "dreamverse")),
_lane("golden-gate", "golden_gate", "lanes/golden_gate.sh"),
_lane("ssim", "ssim", "lanes/ssim.sh", 4, fa4="1"),
_lane("lora-inference", "inference_lora", "lanes/inference_lora.sh"),
_lane("lora-extraction", "lora_extraction", "lanes/lora_extraction.sh"),
_lane("training", "training", "lanes/training.sh", 4),
_lane("distillation", "distillation_dmd", "lanes/distillation_dmd.sh", 2),
_lane("self-forcing", "self_forcing", "lanes/self_forcing.sh", 2),
_lane("lora-training", "training_lora", "lanes/training_lora.sh", 2),
_lane("training-vsa", "training_vsa", "lanes/training_vsa.sh", 2),
_lane("inference-vmoba", "inference_vmoba", "lanes/inference_vmoba.sh"),
_lane("performance", "performance", "lanes/performance.sh", 2),
_lane("api-server", "api_server", "lanes/api_server.sh"),
_lane("train-framework", "train_framework", "lanes/train_framework.sh"),
_lane("eval", "eval", "lanes/eval.sh", extras=("test", "eval-full")),
)
LANES_BY_KEY = {lane["key"]: lane for lane in LANES}
def backend_name(value: str | None, default: str = "slurm") -> str:
backend = value or default
if backend not in BACKENDS:
raise ValueError(f"Unknown CI_GPU_BACKEND: {backend!r}")
return backend
def _identifier(value: str, label: str) -> str:
if not re.fullmatch(r"[A-Za-z0-9][A-Za-z0-9_-]{0,99}", value):
raise ValueError(f"Invalid {label}")
return value
def build_request(env: Mapping[str, str], config: Mapping[str, Any], root: Path) -> dict[str, Any]:
"""Copy only the explicitly accepted Buildkite fields; never forward its token."""
backend = backend_name(env.get("CI_GPU_BACKEND"), config.get("default_backend", "slurm"))
if backend == "slurm":
raise ValueError("Slurm uses the existing trusted uploader and dispatcher")
repository = config.get("repository", DEFAULT_REPOSITORY)
if not re.fullmatch(r"https://github\.com/[A-Za-z0-9_.-]+/[A-Za-z0-9_.-]+\.git", repository):
raise ValueError("repository must be a fixed HTTPS GitHub .git URL")
supplied_repo = env.get("BUILDKITE_REPO", "")
ssh_repo = "git@github.com:" + repository.removeprefix("https://github.com/")
if supplied_repo not in (repository, ssh_repo):
raise ValueError("Buildkite repository does not match operator policy")
commit = env.get("BUILDKITE_COMMIT", "")
if not re.fullmatch(r"[0-9a-f]{40}", commit):
raise ValueError("BUILDKITE_COMMIT must be an immutable 40-character SHA")
build_id = _identifier(env.get("BUILDKITE_BUILD_ID", ""), "BUILDKITE_BUILD_ID")
job_id = _identifier(env.get("BUILDKITE_JOB_ID", ""), "BUILDKITE_JOB_ID")
pr = env.get("BUILDKITE_PULL_REQUEST", "false")
if pr == "false" and env.get("PR_NUMBER"):
pr = env["PR_NUMBER"]
if pr not in ("", "false") and not re.fullmatch(r"[1-9][0-9]*", pr):
raise ValueError("Invalid PR number")
scope = env.get("TEST_SCOPE") or "fastcheck"
if scope not in SCOPES:
raise ValueError(f"Invalid TEST_SCOPE: {scope!r}")
# Reference publication requires a separately trusted publishing job.
if env.get("FASTVIDEO_SSIM_BOOTSTRAP_MODE", "0") not in ("", "0", "false"):
raise ValueError("SSIM bootstrap publication is not supported by the isolated GPU runner")
request = {"build_id": f"{build_id}.{job_id}", "buildkite_build_id": build_id,
"repository": repository, "commit": commit, "pr_number": pr or "false",
"pr_key": f"{repository}#{pr}" if pr not in ("", "false") else None,
"backend": backend, "scope": scope, "env": {"TEST_SCOPE": scope,
"BUILDKITE_PULL_REQUEST": pr or "false", "BUILDKITE_COMMIT": commit,
"BUILDKITE_REPO": repository, "BUILDKITE_BUILD_ID": build_id,
"BUILDKITE_JOB_ID": job_id}}
for key in ("BUILDKITE_BRANCH", "BUILDKITE_SOURCE", "BUILDKITE_BUILD_URL"):
value = env.get(key, "")
if len(value) > 2048 or any(ord(c) < 32 for c in value):
raise ValueError(f"Invalid {key}")
request["env"][key] = value
if scope == "direct":
requested = env.get("TEST_TYPE", "").removesuffix("_ci")
selected = [lane for lane in LANES if requested in (lane["public_type"], lane["key"])]
if len(selected) != 1:
raise ValueError(f"Unknown direct TEST_TYPE: {requested!r}")
elif scope == "merge":
encoded = env.get("MERGE_TEST_PLAN", "")
if not re.fullmatch(r",(?:none|[a-z]+(?:-[a-z]+)*(?:,[a-z]+(?:-[a-z]+)*)*),", encoded):
raise ValueError("Missing or malformed trusted MERGE_TEST_PLAN")
names = encoded[1:-1].split(",")
if names == ["none"]:
names = []
if any(name not in LANES_BY_KEY or LANES_BY_KEY[name]["fastcheck"] for name in names):
raise ValueError("MERGE_TEST_PLAN contains an unknown or non-integration lane")
if len(names) != len(set(names)):
raise ValueError("Duplicate MERGE_TEST_PLAN lane")
selected = [lane for lane in LANES if lane["key"] in names]
elif scope == "scheduled":
requested = env.get("TEST_TYPE", "ssim").removesuffix("_ci")
if requested != "ssim":
raise ValueError("Scheduled scope supports SSIM; performance uses direct scope")
selected = [LANES_BY_KEY[requested]]
else:
selected = [lane for lane in LANES if scope == "full" or lane["fastcheck"]]
for key, lane_key, directory in (("FASTVIDEO_GOLDEN_TEST_FILES", "golden-gate", "golden_gate"),
("FASTVIDEO_SSIM_TEST_FILES", "ssim", "ssim")):
if lane_key not in {lane["key"] for lane in selected}:
continue
merge_key = "MERGE_GOLDEN_TESTS" if lane_key == "golden-gate" else "MERGE_SSIM_TESTS"
value = env.get(merge_key, "") if scope == "merge" else "all"
if not value:
raise ValueError(f"Missing trusted {key} for merge scope")
if value != "all":
names = value.split(",")
if len(names) != len(set(names)) or any(
not re.fullmatch(r"test_[a-z0-9_]+\.py", name) for name in names
):
raise ValueError(f"Invalid or unknown {key}")
# Existence belongs to the isolated exact-SHA checkout. The trusted
# controller need not be redeployed whenever a PR adds a test file.
request["env"][key] = value
request["lanes"] = [{**lane, "script": ".buildkite/scripts/" + lane["script"]} for lane in selected]
return request
+6
View File
@@ -0,0 +1,6 @@
#!/bin/sh
# Install this file at /opt/fastvideo-gpu-ci/run on the trusted dispatch host.
set -eu
exec /opt/fastvideo-gpu-ci/venv/bin/python -I \
/opt/fastvideo-gpu-ci/source/scripts/gpu_ci/entrypoint.py run \
--config /etc/fastvideo-gpu-ci.json
+6
View File
@@ -0,0 +1,6 @@
#!/bin/sh
# Install as the trusted bootstrap for each participating Buildkite pipeline.
set -eu
exec /opt/fastvideo-gpu-ci/venv/bin/python -I \
/opt/fastvideo-gpu-ci/source/scripts/gpu_ci/entrypoint.py upload \
--config /etc/fastvideo-gpu-ci.json
+109
View File
@@ -0,0 +1,109 @@
#!/usr/bin/env bash
# Installed with the trusted controller; never load this entrypoint from a PR.
set -euo pipefail
umask 077
[[ ${FASTVIDEO_CI_REPOSITORY:-} == https://github.com/hao-ai-lab/FastVideo.git ]]
[[ ${FASTVIDEO_CI_COMMIT:-} =~ ^[0-9a-f]{40}$ ]]
[[ ${FASTVIDEO_CI_GPUS:-} =~ ^[1-4]$ ]]
[[ ${FASTVIDEO_CI_EXTRAS:-} =~ ^[a-z][a-z0-9-]*(,[a-z][a-z0-9-]*)*$ ]]
[[ ${FASTVIDEO_CI_SCRIPT:-} == .buildkite/scripts/unit_test.sh ||
${FASTVIDEO_CI_SCRIPT:-} =~ ^\.buildkite/scripts/lanes/[a-z_]+\.sh$ ]]
[[ ${FASTVIDEO_CI_KERNEL:-} =~ ^[01]$ ]]
[[ ${FASTVIDEO_FA4:-} =~ ^[01]$ ]]
export PATH="/opt/venv/bin:/root/.local/bin:$PATH"
export VIRTUAL_ENV=/opt/venv
export HF_HOME=/workspace/hf
export HF_HUB_CACHE="$HF_HOME/hub"
export UV_CACHE_DIR=/workspace/uv-cache
export UV_LINK_MODE=copy
export WANDB_MODE=offline
export FASTVIDEO_CI_LOCAL_ONLY=1
export FASTVIDEO_SSIM_BOOTSTRAP_MODE=0
export MASTER_ADDR=127.0.0.1
export MASTER_PORT=29500
export PYTHONUNBUFFERED=1
mkdir -p /workspace/artifacts "$HF_HUB_CACHE"
# Results are diagnostic only. The controller trusts the container exit status.
finish() {
worker_rc=$?
trap - EXIT
printf '{"exit_code":%s,"commit":"%s"}\n' "$worker_rc" "$FASTVIDEO_CI_COMMIT" \
> /workspace/artifacts/worker-result.json
if [ -d /workspace/repository/fastvideo/tests/ssim/generated_videos ]; then
cp -R /workspace/repository/fastvideo/tests/ssim/generated_videos \
/workspace/artifacts/ssim-generated 2>/dev/null || true
fi
exit "$worker_rc"
}
trap finish EXIT
trap 'exit 143' TERM
trap 'exit 130' INT
# The optional cache is a dedicated, read-only HF Hub subtree. Keep mutable
# refs and locks private; link immutable blobs instead of duplicating weights.
if [ -d /ci-cache ]; then
python - <<'PY'
import os
import shutil
from pathlib import Path
source = Path('/ci-cache').resolve()
destination = Path('/workspace/hf/hub')
for current, directories, filenames in os.walk(source, followlinks=False):
if Path(current) == source:
directories[:] = [name for name in directories if name.startswith(('models--', 'datasets--', 'spaces--'))]
filenames = [] # Never copy a token file from a misconfigured cache root.
directories[:] = [name for name in directories if name != '.locks'
and not (Path(current) / name).is_symlink()]
relative = Path(current).relative_to(source)
target_dir = destination / relative
target_dir.mkdir(parents=True, exist_ok=True)
for name in filenames:
original = Path(current) / name
resolved = original.resolve()
if not resolved.is_relative_to(source) or not resolved.is_file():
raise RuntimeError('Cache contains a file outside its mounted subtree')
target = target_dir / name
if 'blobs' in relative.parts or original.is_symlink():
target.symlink_to(resolved)
else:
shutil.copyfile(original, target)
PY
fi
mkdir /workspace/repository
cd /workspace/repository
git init --quiet
git remote add origin "$FASTVIDEO_CI_REPOSITORY"
if ! git fetch --no-tags --depth=1 origin "$FASTVIDEO_CI_COMMIT"; then
[[ ${FASTVIDEO_CI_PR_NUMBER:-} =~ ^[1-9][0-9]*$ ]]
git fetch --no-tags --depth=1 origin "refs/pull/${FASTVIDEO_CI_PR_NUMBER}/head"
fi
git checkout --quiet --detach FETCH_HEAD
test "$(git rev-parse HEAD)" = "$FASTVIDEO_CI_COMMIT"
git -c protocol.file.allow=never submodule update --init --recursive --depth=1
# Installing dependencies after the kernel can silently replace the PR kernel.
printf 'fastvideo-kernel\n' > /workspace/kernel-excludes
uv pip install --excludes /workspace/kernel-excludes -e ".[${FASTVIDEO_CI_EXTRAS}]"
if [ "$FASTVIDEO_CI_KERNEL" = 1 ]; then
python fastvideo/tests/modal/kernel_build_cache.py install
fi
python - <<'PY'
import os
import torch
expected = int(os.environ['FASTVIDEO_CI_GPUS'])
if not torch.cuda.is_available() or torch.cuda.device_count() != expected:
raise RuntimeError(f'Expected exactly {expected} CUDA devices')
for index in range(expected):
print(f'CUDA {index}: {torch.cuda.get_device_name(index)}', flush=True)
PY
# SSIM runs concurrent pytest subprocesses; a single global JUnit destination
# would race. Preserve each lane's native metrics plus the controller logs.
bash "$FASTVIDEO_CI_SCRIPT"