Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
14 changes: 14 additions & 0 deletions .github/workflows/e2e-gpu-job.yml
Original file line number Diff line number Diff line change
Expand Up @@ -39,6 +39,11 @@ on:
type: string
default: ""
description: "pytest -k filter expression"
extra_models:
required: false
type: string
default: ""
description: "Extra model ids to download by id (space-separated), for skip_tier_download models"
setup_agentic_deps:
required: false
type: boolean
Expand Down Expand Up @@ -126,6 +131,15 @@ jobs:
- name: Download models
run: bash scripts/ci_download_model.sh --gpu-tier ${{ inputs.gpu_tier }}

- name: Download extra models
if: inputs.extra_models != ''
# Pass the input via env (not interpolated into the run script) so its
# text can never become shell syntax; unquoted $EXTRA_MODELS still
# word-splits into multiple space-separated model ids.
env:
EXTRA_MODELS: ${{ inputs.extra_models }}
run: bash scripts/ci_download_model.sh $EXTRA_MODELS
Comment thread
slin1237 marked this conversation as resolved.

# Run tests
- name: Run E2E tests
timeout-minutes: ${{ inputs.test_timeout }}
Expand Down
17 changes: 16 additions & 1 deletion .github/workflows/pr-test-rust.yml
Original file line number Diff line number Diff line change
Expand Up @@ -732,6 +732,20 @@ jobs:
test_filter: "-k TestIGWMixedWorkerClassification"
secrets: inherit

e2e-4gpu-epd:
name: e2e-4gpu-epd (tokenspeed)
needs: [e2e-1gpu-chat]
Comment on lines +735 to +737

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge Include EPD job in the finish gate

Adding this workflow job does not make its result part of the aggregate finish check: I inspected the finish job in this same workflow, and its needs list and failure condition still omit e2e-4gpu-epd. If branch protection relies on finish, this new smoke can fail or still be running while finish reports success, so the EPD coverage added here would not actually block merges.

Useful? React with 👍 / 👎.

uses: ./.github/workflows/e2e-gpu-job.yml
with:
engine: tokenspeed
gpu_tier: "4"
runner: 4-gpu-h100
timeout: 75
test_timeout: 65
test_dirs: e2e_test/chat_completions/test_epd_multimodal.py
extra_models: "Qwen/Qwen3.5-9B"
secrets: inherit

# --- Vendor E2E: CPU-only cloud backend tests ---

e2e-vendor:
Expand Down Expand Up @@ -1023,7 +1037,7 @@ jobs:
path: benchmark_go_bindings/

finish:
needs: [pre-commit, python-lint, grpc-proto-build-check, build-wheel, python-unit-tests, unit-tests, benchmarks, e2e-1gpu-chat, e2e-1gpu-completions, e2e-1gpu-embeddings, e2e-1gpu-gateway, e2e-1gpu-responses, e2e-2gpu-pd, e2e-4gpu-chat, e2e-4gpu-gateway, e2e-vendor, go-unit-tests, go-bindings-e2e]
needs: [pre-commit, python-lint, grpc-proto-build-check, build-wheel, python-unit-tests, unit-tests, benchmarks, e2e-1gpu-chat, e2e-1gpu-completions, e2e-1gpu-embeddings, e2e-1gpu-gateway, e2e-1gpu-responses, e2e-2gpu-pd, e2e-4gpu-chat, e2e-4gpu-gateway, e2e-4gpu-epd, e2e-vendor, go-unit-tests, go-bindings-e2e]
if: always()
runs-on: k8s-runner-cpu
permissions: {}
Expand All @@ -1045,6 +1059,7 @@ jobs:
"${{ needs.e2e-2gpu-pd.result }}" == "failure" || \
"${{ needs.e2e-4gpu-chat.result }}" == "failure" || \
"${{ needs.e2e-4gpu-gateway.result }}" == "failure" || \
"${{ needs.e2e-4gpu-epd.result }}" == "failure" || \
"${{ needs.e2e-vendor.result }}" == "failure" || \
"${{ needs.go-unit-tests.result }}" == "failure" || \
"${{ needs.go-bindings-e2e.result }}" == "failure" ]]; then
Expand Down
151 changes: 151 additions & 0 deletions e2e_test/chat_completions/test_epd_multimodal.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,151 @@
"""EPD (Encode-Prefill-Decode) multimodal Chat Completions E2E tests.

Exercises TokenSpeed's EPD disaggregation on a vision-language model across four
worker-count topologies: the encode worker runs the vision tower, prefill/decode
run the LM, and the gateway stitches encode -> prefill -> decode.

The gateway runs EPD-only (``RoutingMode::EncodePrefillDecode``) — there is no
single-worker fallback path — so a *correct* answer about the image proves the
encode->prefill->decode pipeline ran for that request: the model cannot name the
color/animal unless the encoder's embeddings reached prefill+decode. As an extra
per-request check we assert the router logged a fresh EPD encode dispatch. Qwen3.5
is a thinking model, so the answer may arrive in ``reasoning_content`` rather than
``content`` — we accept either channel.

Usage:
pytest e2e_test/chat_completions/test_epd_multimodal.py -v
"""

from __future__ import annotations

import base64
import logging
import tempfile
import time
from pathlib import Path

import pytest
from infra.pd_logs import LOG_FLUSH_TIMEOUT_S, read_logs

logger = logging.getLogger(__name__)

FIXTURES_DIR = Path(__file__).parent.parent / "fixtures" / "images"
DOG_IMAGE_PATH = FIXTURES_DIR / "dog.jpg" # Black labrador puppy (checked in)
PUG_IMAGE_PATH = FIXTURES_DIR / "pug.jpg" # Pug in a blanket (checked in)

# Solid 32x32 RGB PNGs (from docs/guides/epd-ts-test.md). A solid color is the
# most deterministic vision check: the model can only name it if the encoder's
# pixels actually reached the LM through prefill+decode.
RED_PNG_B64 = (
"iVBORw0KGgoAAAANSUhEUgAAACAAAAAgCAIAAAD8GO2jAAAAKElEQVR4nO3NsQ0AAAzCMP5/"
"un0CNkuZ41wybXsHAAAAAAAAAAAAxR4yw/wuPL6QkAAAAABJRU5ErkJggg=="
)
BLUE_PNG_B64 = (
"iVBORw0KGgoAAAANSUhEUgAAACAAAAAgCAIAAAD8GO2jAAAAJklEQVR4nO3NsQkAAAjAsP7/"
"tF7hIASyp5pjAoFAIBAIBAKB4EmwOkv8Lm7+zY4AAAAASUVORK5CYII="
)

_LOG_DIR = Path(tempfile.mkdtemp(prefix="smg-e2e-epd-"))
Comment thread
slin1237 marked this conversation as resolved.
# Per-request, router-side proof that the gateway routed through the encode stage.
# The worker-side "EPD encode: accepted" INFO line does not survive TokenSpeed's
# logging reconfiguration; this router line does, and the gateway runs at
# log_level=debug (below).
EPD_DISPATCH_MARKER = "EPD encode dispatch issued"

# (encode, prefill, decode) worker counts. Every worker is tp=1, so 1e1p1d uses
# 3 GPUs and the rest use 4 — all fit the 4-GPU runner. Counts ride in the param
# because setup_backend is class-scoped and can't read per-param marks.
_EPD_TOPOLOGIES = [
pytest.param(("epd_grpc", (1, 1, 1)), id="1e1p1d"),
pytest.param(("epd_grpc", (1, 2, 1)), id="1e2p1d"),
pytest.param(("epd_grpc", (2, 1, 1)), id="2e1p1d"),
pytest.param(("epd_grpc", (1, 1, 2)), id="1e1p2d"),
]


def _file_to_data_url(path: Path) -> str:
data = base64.b64encode(path.read_bytes()).decode("utf-8")
return f"data:image/jpeg;base64,{data}"


def _b64_png_to_data_url(data_b64: str) -> str:
return f"data:image/png;base64,{data_b64}"


def _epd_dispatch_count() -> int:
"""How many EPD encode dispatches the router has logged so far (cumulative)."""
return read_logs(_LOG_DIR, "smg*").count(EPD_DISPATCH_MARKER)


@pytest.mark.engine("tokenspeed")
@pytest.mark.gpu(4)
@pytest.mark.e2e
@pytest.mark.model("Qwen/Qwen3.5-9B")
@pytest.mark.gateway(log_level="debug", policy="cache_aware", log_dir=str(_LOG_DIR))
@pytest.mark.parametrize("setup_backend", _EPD_TOPOLOGIES, indirect=True)
class TestEPDMultimodal:
"""Verify the image really flows encode -> prefill -> decode for each topology."""

def _check_image(self, client, model, image_url, question, keywords):
"""Send one image request; assert a correct answer + a fresh EPD dispatch."""
# Baseline BEFORE the request: the router log accumulates across topologies
# and is never cleared, so assert THIS request added a new dispatch.
dispatches_before = _epd_dispatch_count()

response = client.chat.completions.create(
model=model,
messages=[
{
"role": "user",
"content": [
{"type": "text", "text": question},
{"type": "image_url", "image_url": {"url": image_url}},
],
}
],
temperature=0,
max_tokens=256,
)

# (1) The EPD pipeline produced output. A thinking model may put the answer
# in reasoning_content instead of content — accept either. Dump the whole
# message on failure so we can see where (if anywhere) the answer landed.
msg = response.choices[0].message.model_dump()
content = msg.get("content") or ""
reasoning = msg.get("reasoning_content") or msg.get("reasoning") or ""
text = f"{content}\n{reasoning}".strip()
assert text, f"EPD pipeline returned empty content AND reasoning; message={msg}"

# (2) The answer is correct about the image — only possible if the encoder's
# embeddings reached prefill+decode (EPD-only gateway; no fallback path).
assert any(k in text.lower() for k in keywords), (
f"expected one of {keywords} in the answer, got: {text!r}"
)
Comment on lines +119 to +123

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🎯 Functional Correctness | 🟠 Major | ⚡ Quick win

Require answers that distinguish each image.

Substring matching allows false passes: both animal cases accept "dog"/"puppy", and generic or negated text containing a color also passes. The test can therefore succeed without proving that different images produced different classifications.

Use normalized exact answers for colors and disjoint expectations—such as requiring "pug" for the pug image—or use images from different species.

Also applies to: 136-152

🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@e2e_test/chat_completions/test_epd_multimodal.py` around lines 120 - 124,
Replace the broad substring assertions in the multimodal image tests around the
answer validation with normalized exact-answer checks and disjoint expected
classifications for each image. Require the pug case to return “pug” and ensure
color cases accept only their intended normalized color, preventing generic or
negated text from passing.


# (3) The router logged a fresh EPD encode dispatch for THIS request.
deadline = time.monotonic() + LOG_FLUSH_TIMEOUT_S
while _epd_dispatch_count() <= dispatches_before and time.monotonic() < deadline:
time.sleep(0.5)
assert _epd_dispatch_count() > dispatches_before, (
"router logged no new EPD encode dispatch for this request; the gateway "
f"did not route through encode->prefill->decode (checked {_LOG_DIR}/smg*)"
)
logger.info("EPD OK (%s): %s", keywords[0], text[:120])

def test_color_images(self, model, setup_backend):
"""Solid red then blue — the most deterministic encode->decode check."""
_, _, client, *_ = setup_backend
question = "What color is this image? Reply with just the color."
self._check_image(client, model, _b64_png_to_data_url(RED_PNG_B64), question, ["red"])
self._check_image(client, model, _b64_png_to_data_url(BLUE_PNG_B64), question, ["blue"])

def test_animal_images(self, model, setup_backend):
"""Real photos — dog then pug."""
_, _, client, *_ = setup_backend
question = "What animal is in this image?"
self._check_image(
client, model, _file_to_data_url(DOG_IMAGE_PATH), question, ["dog", "puppy", "labrador"]
)
self._check_image(
client, model, _file_to_data_url(PUG_IMAGE_PATH), question, ["pug", "dog", "puppy"]
)
110 changes: 107 additions & 3 deletions e2e_test/fixtures/setup_backend.py
Original file line number Diff line number Diff line change
Expand Up @@ -120,7 +120,11 @@ def setup_backend(request: pytest.FixtureRequest):
Returns:
Tuple of ``(backend_name, model_path, client, gateway)``
"""
backend_name: str = request.param
raw_param = request.param
if isinstance(raw_param, tuple):
backend_name, epd_counts = raw_param
else:
backend_name, epd_counts = raw_param, None

if os.environ.get(ENV_SKIP_BACKEND_SETUP, "").lower() in ("1", "true", "yes"):
pytest.skip(f"{ENV_SKIP_BACKEND_SETUP} is set")
Expand All @@ -137,8 +141,9 @@ def setup_backend(request: pytest.FixtureRequest):
return

# Local backends
is_epd = backend_name.startswith("epd_")
is_pd = backend_name.startswith("pd_")
protocol = backend_name.replace("pd_", "")
protocol = backend_name.replace("epd_", "").replace("pd_", "")
connection_mode = ConnectionMode(protocol)
engine = get_runtime()
model_path = get_model_spec(model_id)["model"]
Expand All @@ -154,7 +159,18 @@ def setup_backend(request: pytest.FixtureRequest):

gateway = Gateway()
try:
if is_pd:
if is_epd:
yield from _setup_epd(
model_id,
model_path,
engine,
connection_mode,
epd_counts,
gateway_config,
gateway,
log_dir,
)
elif is_pd:
yield from _setup_pd(
model_id,
model_path,
Expand Down Expand Up @@ -299,6 +315,94 @@ def _setup_pd(
stop_workers(all_workers)


# ---------------------------------------------------------------------------
# EPD (encode-prefill-decode) disaggregation backend
# ---------------------------------------------------------------------------


def _setup_epd(
model_id,
model_path,
engine,
connection_mode,
epd_counts,
gateway_config,
gateway,
log_dir,
):
"""Launch encode + prefill + decode workers + EPD gateway, yield, tear down.

``epd_counts`` is ``(n_encode, n_prefill, n_decode)``. Every worker is tp=1
(one GPU); GPUs are assigned sequentially E -> P -> D.
"""
if not epd_counts or len(epd_counts) != 3:
raise ValueError("epd_grpc backend requires a (n_encode, n_prefill, n_decode) param")
n_encode, n_prefill, n_decode = epd_counts
spec = get_model_spec(model_id)
tp = spec.get("tp", 1)
backend_name = f"epd_{connection_mode.value}"
runtime_label = RUNTIME_LABELS.get(engine, engine)

logger.info(
"Starting %s EPD backend: model=%s, %de + %dp + %dd",
runtime_label,
model_id,
n_encode,
n_prefill,
n_decode,
)

all_workers: list = []
try:
encode_workers = _start_workers_tracked(
model_id=model_id,
engine=engine,
mode=connection_mode,
count=n_encode,
worker_type=WorkerType.ENCODE,
log_dir=log_dir,
gpu_offset=0,
gpus=1, # vision tower runs on one GPU regardless of LM tp
)
all_workers.extend(encode_workers)

prefill_workers = _start_workers_tracked(
model_id=model_id,
engine=engine,
mode=connection_mode,
count=n_prefill,
worker_type=WorkerType.PREFILL,
log_dir=log_dir,
gpu_offset=n_encode,
)
all_workers.extend(prefill_workers)

decode_workers = _start_workers_tracked(
model_id=model_id,
engine=engine,
mode=connection_mode,
count=n_decode,
worker_type=WorkerType.DECODE,
log_dir=log_dir,
gpu_offset=n_encode + n_prefill * tp,
)
all_workers.extend(decode_workers)

_start_gateway(
gateway,
gateway_config,
encode_workers=encode_workers,
prefill_workers=prefill_workers,
decode_workers=decode_workers,
)
logger.info("%s EPD backend ready at %s", runtime_label, gateway.base_url)
yield backend_name, model_path, _make_openai_client(gateway), gateway
finally:
logger.info("Tearing down %s EPD backend", runtime_label)
gateway.shutdown()
stop_workers(all_workers)


# ---------------------------------------------------------------------------
# Cloud backend
# ---------------------------------------------------------------------------
Expand Down
1 change: 1 addition & 0 deletions e2e_test/infra/constants.py
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@ class WorkerType(StrEnum):
"""Worker specialization type."""

REGULAR = "regular"
ENCODE = "encode"
PREFILL = "prefill"
DECODE = "decode"

Expand Down
Loading
Loading