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
2 changes: 2 additions & 0 deletions example-expb.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -86,6 +86,8 @@ scenarios:
network: mainnet
# Optional: Snapshot backend to use. Available snapshot backends: overlay, zfs, copy. Defaults to overlay
snapshot_backend: overlay
# Optional: Snapshot mount target. Use /execution-data/geth when the source contains the contents of Geth's datadir/geth.
# snapshot_mount_path: /execution-data/geth
# Optional: Snapshot path for copy backend. If defined, this path will be used instead of work_dir/snapshot. Only used when snapshot_backend is copy
snapshot_path: ./work_dir/snapshot
# Optional: Number of times to repeat the scenario. Defaults to 1.
Expand Down
5 changes: 5 additions & 0 deletions src/expb/configs/scenarios.py
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
from pathlib import Path
from typing import Literal

from pydantic import (
BaseModel,
Expand Down Expand Up @@ -139,6 +140,10 @@ class Scenario(BaseModel):
description="Snapshot backend to use.",
default=SnapshotBackend.OVERLAY,
)
snapshot_mount_path: Literal["/execution-data", "/execution-data/geth"] = Field(
description="Path where the snapshot is mounted in the execution client container.",
default="/execution-data",
)
snapshot_path: Path | None = Field(
description="Path to the snapshot directory for copy backend (overrides work_dir).",
default=None,
Expand Down
4 changes: 2 additions & 2 deletions src/expb/payloads/executor/executor_config.py
Original file line number Diff line number Diff line change
Expand Up @@ -14,7 +14,6 @@
CLIENT_METRICS_PORT,
CLIENT_RPC_PORT,
CLIENT_RPC_WS_PORT,
CLIENTS_DATA_DIR,
CLIENTS_JWT_SECRET_DIR,
Client,
)
Expand Down Expand Up @@ -124,6 +123,7 @@ def __init__(

## Snapshot config
self.snapshot_source: str = scenario.snapshot_source
self.snapshot_mount_path: str = scenario.snapshot_mount_path
self.snapshot_service: SnapshotService = snapshot_service

## Outputs directory
Expand Down Expand Up @@ -315,7 +315,7 @@ def get_execution_client_volumes(self) -> list[dict[str, Any]]:
)
execution_container_volumes.append(
{
"bind": CLIENTS_DATA_DIR,
"bind": self.snapshot_mount_path,
"config": {
"name": f"{container_name}-overlay-merged",
"driver": "local",
Expand Down
2 changes: 2 additions & 0 deletions src/expb/payloads/executor/services/k6.py
Original file line number Diff line number Diff line change
Expand Up @@ -52,6 +52,8 @@ def build_k6_script_config(
"thresholds": {
"http_req_failed{kind:newPayload}": ["rate < 0.01"],
"http_req_failed{kind:forkchoiceUpdated}": ["rate < 0.01"],
"checks{kind:newPayload}": ["rate == 1"],
"checks{kind:forkchoiceUpdated}": ["rate == 1"],
},
"systemTags": [
"scenario",
Expand Down
66 changes: 61 additions & 5 deletions src/expb/payloads/executor/services/templates/k6-script.js.j2
Original file line number Diff line number Diff line change
Expand Up @@ -15,10 +15,10 @@ const abortOnEOF = (__ENV.EXPB_ABORT_ON_EOF || '0') === '1';
const addCorrelationHeader = (__ENV.EXPB_ADD_CID || '1') === '1';
const engineEndpoint = __ENV.EXPB_ENGINE_ENDPOINT;
const perPayloadMetrics = (__ENV.EXPB_PER_PAYLOAD_METRICS || '0') === '1';
const discardResponses = (__ENV.EXPB_DISCARD_RESPONSES || '1') === '1';
const enableLogging = (__ENV.EXPB_ENABLE_LOGGING || '0') === '1';
const perPayloadMetricsLogs = (__ENV.EXPB_PER_PAYLOAD_METRICS_LOGS || '0') === '1';
const verbosePostLogs = (__ENV.EXPB_VERBOSE_POST_LOGS || '0') === '1';
const maxDiagnosticBodyLength = 512;

// Load k6 options JSON
const config = JSON.parse(open(__ENV.EXPB_CONFIG_FILE_PATH));
Expand Down Expand Up @@ -76,6 +76,62 @@ function emitPerPayloadMetricsRow(idx, gasUsed, processingMs) {
);
}

function boundedResponseBody(response) {
const body = typeof response.body === 'string' ? response.body : '';
const compactBody = body.replace(/\s+/g, ' ');
if (compactBody.length <= maxDiagnosticBodyLength) return compactBody;
return compactBody.slice(0, maxDiagnosticBodyLength) + '...';
}

function checkEngineResponse(response, kind, idx, tags) {
let payloadStatus = null;
let valid = response.status === 200;
let diagnostic = response.status === 200 ? '' : `HTTP status ${response.status}`;

if (valid) {
let body;
try {
body = JSON.parse(response.body || '');
} catch (error) {
valid = false;
diagnostic = 'response body is not valid JSON';
}

if (valid && (!body || body.error !== undefined)) {
valid = false;
diagnostic = 'JSON-RPC error';
}

if (valid && (!body || body.result === undefined)) {
valid = false;
diagnostic = 'response has no JSON-RPC result';
}

if (valid) {
// engine_newPayload returns the payloadStatus object directly, while
// engine_forkchoiceUpdated wraps it in result.payloadStatus.
payloadStatus = kind === 'newPayload'
? body.result && body.result.status
: body.result && body.result.payloadStatus && body.result.payloadStatus.status;
if (payloadStatus !== 'VALID') {
valid = false;
diagnostic = `payload status ${payloadStatus || 'missing'}`;
}
}
}

if (!valid) {
console.error(
`[engine] ${kind} idx=${idx} ${diagnostic}; response=${boundedResponseBody(response)}`
);
}

check(response, {
'status_200': (x) => x.status === 200,
'payload_status_valid': () => valid,
}, tags);
}

// --- Fetch next pair from payload server ---
// Line format: {metadata_json}\t{raw_NP}\t{raw_FCU}
// Metadata is pre-extracted by the executor, no regex needed here.
Expand Down Expand Up @@ -170,9 +226,9 @@ function processNextPayload(base_tags = {}, warmup = false) {
const r = http.post(engineEndpoint, pair.rawPayload, {
headers: headers,
tags: tags,
responseType: discardResponses ? 'none' : 'text',
responseType: 'text',
});
check(r, { 'status_200': (x) => x.status === 200 }, tags);
checkEngineResponse(r, 'newPayload', idx, tags);

if (!warmup) {
emitPerPayloadMetricsRow(idx, pair.gasUsed, r.timings.waiting);
Expand Down Expand Up @@ -213,9 +269,9 @@ function processNextPayload(base_tags = {}, warmup = false) {
const r = http.post(engineEndpoint, pair.rawFcu, {
headers: headers,
tags: tags,
responseType: discardResponses ? 'none' : 'text',
responseType: 'text',
});
check(r, { 'status_200': (x) => x.status === 200 }, tags);
checkEngineResponse(r, 'forkchoiceUpdated', idx, tags);
});

// Emit client-side processing metric for the PREVIOUS block.
Expand Down
89 changes: 89 additions & 0 deletions tests/test_executor_config.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,89 @@
from pathlib import Path

import pytest

from expb.clients import Client
from expb.configs.scenarios import Scenario, ScenariosPaths
from expb.payloads.executor import executor_config as executor_config_module
from expb.payloads.executor.executor_config import ExecutorConfig
from expb.payloads.executor.services.snapshots import SnapshotService


class StubSnapshotService(SnapshotService):
def __init__(self, snapshot_path: Path):
self.snapshot_path = snapshot_path

def get_snapshot(self, name: str, source: str) -> Path:
return self.snapshot_path


def make_scenario(tmp_path: Path, **overrides) -> Scenario:
payloads = tmp_path / "payloads.jsonl"
fcus = tmp_path / "fcus.jsonl"
payloads.touch()
fcus.touch()
values = {
"client": Client.GETH,
"payloads": payloads,
"fcus": fcus,
"snapshot_source": "snapshot",
}
values.update(overrides)
return Scenario.model_validate(values)


@pytest.mark.parametrize(
"mount_path",
["/execution-data", "/execution-data/geth"],
)
def test_snapshot_mount_path_allows_supported_client_layouts(
tmp_path: Path, mount_path: str
) -> None:
scenario = make_scenario(tmp_path, snapshot_mount_path=mount_path)

assert scenario.snapshot_mount_path == mount_path


def test_snapshot_mount_path_defaults_to_client_data_directory(tmp_path: Path) -> None:
scenario = make_scenario(tmp_path)

assert scenario.snapshot_mount_path == "/execution-data"


def test_snapshot_mount_path_rejects_arbitrary_container_paths(tmp_path: Path) -> None:
with pytest.raises(ValueError):
make_scenario(tmp_path, snapshot_mount_path="/tmp/snapshot")


def test_executor_mounts_snapshot_at_configured_path(
tmp_path: Path, monkeypatch: pytest.MonkeyPatch
) -> None:
monkeypatch.setattr(executor_config_module.docker, "from_env", lambda: object())
monkeypatch.setattr(executor_config_module.os, "getuid", lambda: 0, raising=False)
monkeypatch.setattr(executor_config_module.os, "getgid", lambda: 0, raising=False)

snapshot_path = tmp_path / "snapshot"
snapshot_path.mkdir()
scenario = make_scenario(tmp_path, name="geth-fusaka", snapshot_mount_path="/execution-data/geth")
config = ExecutorConfig(
scenario=scenario,
snapshot_service=StubSnapshotService(snapshot_path),
paths=ScenariosPaths(
work=tmp_path / "work",
outputs=tmp_path / "outputs",
),
)

volumes = config.get_execution_client_volumes()
snapshot_volume = next(
volume
for volume in volumes
if volume["config"]["name"].endswith("-overlay-merged")
)

assert snapshot_volume["bind"] == "/execution-data/geth"
assert snapshot_volume["config"]["driver_opts"]["o"].startswith("bind,rw")
assert snapshot_volume["config"]["driver_opts"]["device"] == str(
snapshot_path.resolve()
)
assert all(volume["bind"] != "/execution-data" for volume in volumes)
95 changes: 95 additions & 0 deletions tests/test_k6.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,95 @@
import json
import re
import shutil
import subprocess

import pytest

from expb.clients import Client
from expb.payloads.executor.services.k6 import (
build_k6_script_config,
get_k6_script_content,
)


def test_k6_thresholds_fail_non_valid_engine_responses() -> None:
config = build_k6_script_config(
test_id="test",
scenario_name="scenario",
client=Client.GETH,
iterations=1,
)

thresholds = config["options"]["thresholds"]
assert thresholds["checks{kind:newPayload}"] == ["rate == 1"]
assert thresholds["checks{kind:forkchoiceUpdated}"] == ["rate == 1"]


def test_k6_engine_response_validation_covers_engine_statuses() -> None:
node = shutil.which("node")
if node is None:
pytest.skip("node is required to execute the K6 response validation")

script = get_k6_script_content()
script = re.sub(r"^import .*;\r?\n", "", script, flags=re.MULTILINE)
script = script.replace("export const options", "const options")
script = script.replace("export function setup", "function setup")
script = script.replace("export default function ()", "function defaultFunction()")
script += "\nglobalThis.checkEngineResponse = checkEngineResponse;\n"

driver = f"""
const vm = require('node:vm');
const checks = [];
const errors = [];
const context = {{
__ENV: {{
EXPB_CONFIG_FILE_PATH: '/config.json',
EXPB_JWTSECRET_FILE_PATH: '/jwt.hex',
}},
open: (path) => path === '/config.json' ? '{{"options":{{}}}}' : '00'.repeat(32),
Gauge: function Gauge() {{ this.add = () => {{}}; }},
encoding: {{ b64encode: () => '' }},
check: (response, checkFunctions, tags) => {{
const result = {{}};
for (const [name, checkFunction] of Object.entries(checkFunctions)) {{
result[name] = Boolean(checkFunction(response));
}}
checks.push({{ result, tags }});
return Object.values(result).every(Boolean);
}},
console: {{ error: (message) => errors.push(message) }},
}};
vm.runInNewContext({json.dumps(script)}, context);
const cases = [
['newPayload', {{ status: 200, body: '{{"result":{{"status":"VALID"}}}}' }}, true],
['forkchoiceUpdated', {{ status: 200, body: '{{"result":{{"payloadStatus":{{"status":"VALID"}}}}}}' }}, true],
['newPayload', {{ status: 200, body: '{{"result":{{"status":"INVALID"}}}}' }}, false],
['newPayload', {{ status: 200, body: '{{"result":{{"status":"SYNCING"}}}}' }}, false],
['newPayload', {{ status: 200, body: '{{"result":{{"status":"ACCEPTED"}}}}' }}, false],
['forkchoiceUpdated', {{ status: 200, body: '{{"result":{{"payloadStatus":{{"status":"INVALID"}}}}}}' }}, false],
['forkchoiceUpdated', {{ status: 200, body: '{{"result":{{"payloadStatus":{{"status":"SYNCING"}}}}}}' }}, false],
['forkchoiceUpdated', {{ status: 200, body: '{{"result":{{"payloadStatus":{{"status":"ACCEPTED"}}}}}}' }}, false],
['newPayload', {{ status: 200, body: '{{"error":{{"code":-32000,"message":"failed"}}}}' }}, false],
['forkchoiceUpdated', {{ status: 200, body: '{{' }}, false],
['newPayload', {{ status: 503, body: 'temporarily unavailable' }}, false],
['forkchoiceUpdated', {{ status: 200, body: 'x'.repeat(2000) }}, false],
];
for (const [kind, response, expected] of cases) {{
context.checkEngineResponse(response, kind, 42, {{ kind }});
if (checks.at(-1).result.payload_status_valid !== expected) {{
throw new Error(`unexpected result for ${{kind}}: ${{JSON.stringify(response)}}`);
}}
}}
if (errors.length !== cases.length - 2) throw new Error(`unexpected diagnostic count: ${{errors.length}}`);
if (errors.some((message) => message.length > 700)) throw new Error('diagnostic was not bounded');
if (!errors.some((message) => message.includes('...'))) throw new Error('long diagnostic was not truncated');
process.stdout.write(JSON.stringify({{ checks, errors }}));
"""
result = subprocess.run(
[node, "-e", driver],
capture_output=True,
text=True,
)

assert result.returncode == 0, result.stderr
assert len(json.loads(result.stdout)["checks"]) == 12
Loading