From 4953272bcfbb8055531905811b6c65def6f7583e Mon Sep 17 00:00:00 2001 From: "a.tolmachev" Date: Mon, 6 Apr 2026 12:12:21 +0300 Subject: [PATCH] feat: scaffold streaming e2e fixtures --- TASKS.md | 2 +- apps/ui/scripts/playwright-stack.sh | 51 ++++++ apps/ui/scripts/playwright-ui-server.js | 14 ++ apps/ui/scripts/stream-fixture-server.js | 145 ++++++++++++++++++ .../examples/stream_fixture.rs | 22 +++ 5 files changed, 233 insertions(+), 1 deletion(-) create mode 100644 apps/ui/scripts/stream-fixture-server.js create mode 100644 crates/crank-adapter-grpc/examples/stream_fixture.rs diff --git a/TASKS.md b/TASKS.md index f19acb4..64c6dfb 100644 --- a/TASKS.md +++ b/TASKS.md @@ -4,7 +4,7 @@ ### `feat/streaming-e2e` -Status: pending +Status: in progress DoD: - local fixture stack includes: diff --git a/apps/ui/scripts/playwright-stack.sh b/apps/ui/scripts/playwright-stack.sh index f5ff65d..13589a4 100644 --- a/apps/ui/scripts/playwright-stack.sh +++ b/apps/ui/scripts/playwright-stack.sh @@ -9,6 +9,8 @@ POSTGRES_PORT="${CRANK_E2E_POSTGRES_PORT:-55433}" ADMIN_PORT="${CRANK_E2E_ADMIN_PORT:-3301}" MCP_PORT="${CRANK_E2E_MCP_PORT:-3302}" UI_PORT="${CRANK_E2E_UI_PORT:-3300}" +STREAM_FIXTURE_PORT="${CRANK_E2E_STREAM_FIXTURE_PORT:-3310}" +GRPC_FIXTURE_BIND="${CRANK_E2E_GRPC_FIXTURE_BIND:-127.0.0.1:3311}" POSTGRES_DB="${CRANK_E2E_POSTGRES_DB:-crank}" POSTGRES_USER="${CRANK_E2E_POSTGRES_USER:-crank}" POSTGRES_PASSWORD="${CRANK_E2E_POSTGRES_PASSWORD:-crank}" @@ -38,16 +40,45 @@ cleanup() { kill "$(cat "$TMP_DIR/ui-server.pid")" >/dev/null 2>&1 || true rm -f "$TMP_DIR/ui-server.pid" fi + if [[ -f "$TMP_DIR/stream-fixture.pid" ]]; then + kill "$(cat "$TMP_DIR/stream-fixture.pid")" >/dev/null 2>&1 || true + rm -f "$TMP_DIR/stream-fixture.pid" + fi + if [[ -f "$TMP_DIR/grpc-fixture.pid" ]]; then + kill "$(cat "$TMP_DIR/grpc-fixture.pid")" >/dev/null 2>&1 || true + rm -f "$TMP_DIR/grpc-fixture.pid" + fi docker rm -f "$POSTGRES_CONTAINER" >/dev/null 2>&1 || true kill_port_processes "$UI_PORT" kill_port_processes "$ADMIN_PORT" kill_port_processes "$MCP_PORT" + kill_port_processes "$STREAM_FIXTURE_PORT" + kill_port_processes "${GRPC_FIXTURE_BIND##*:}" } trap cleanup EXIT INT TERM cleanup +wait_for_port() { + local host="$1" + local port="$2" + until python - "$host" "$port" <<'PY' +import socket, sys +sock = socket.socket() +sock.settimeout(0.5) +try: + sock.connect((sys.argv[1], int(sys.argv[2]))) +except OSError: + sys.exit(1) +finally: + sock.close() +PY + do + sleep 1 + done +} + docker run -d --rm \ --name "$POSTGRES_CONTAINER" \ -e POSTGRES_DB="$POSTGRES_DB" \ @@ -76,6 +107,8 @@ export CRANK_BOOTSTRAP_ADMIN_DISPLAY_NAME="Crank E2E" export CRANK_DEMO_SEED="true" export CRANK_PUBLIC_BASE_URL="http://127.0.0.1:$UI_PORT" export CRANK_MCP_PUBLIC_URL="http://127.0.0.1:$MCP_PORT" +export CRANK_E2E_STREAM_FIXTURE_PORT="$STREAM_FIXTURE_PORT" +export CRANK_E2E_GRPC_FIXTURE_BIND="$GRPC_FIXTURE_BIND" mkdir -p "$CRANK_STORAGE_ROOT" @@ -99,6 +132,24 @@ until curl -fsS "http://127.0.0.1:$MCP_PORT/health" >/dev/null 2>&1; do sleep 1 done +( + cd "$ROOT_DIR/apps/ui" + node scripts/stream-fixture-server.js >"$LOG_DIR/stream-fixture.log" 2>&1 +) & +echo $! > "$TMP_DIR/stream-fixture.pid" + +until curl -fsS "http://127.0.0.1:$STREAM_FIXTURE_PORT/health" >/dev/null 2>&1; do + sleep 1 +done + +( + cd "$ROOT_DIR" + cargo run -p crank-adapter-grpc --features test-support --example stream_fixture >"$LOG_DIR/grpc-fixture.log" 2>&1 +) & +echo $! > "$TMP_DIR/grpc-fixture.pid" + +wait_for_port "${GRPC_FIXTURE_BIND%%:*}" "${GRPC_FIXTURE_BIND##*:}" + ( cd "$ROOT_DIR/apps/ui" node scripts/playwright-ui-server.js >"$LOG_DIR/ui-server.log" 2>&1 diff --git a/apps/ui/scripts/playwright-ui-server.js b/apps/ui/scripts/playwright-ui-server.js index 2e1f8ef..9f03f4b 100644 --- a/apps/ui/scripts/playwright-ui-server.js +++ b/apps/ui/scripts/playwright-ui-server.js @@ -50,6 +50,12 @@ function mapRoute(urlPath) { if (urlPath === '/usage') { return path.join(ROOT_DIR, 'html', 'usage.html'); } + if (urlPath === '/stream-sessions') { + return path.join(ROOT_DIR, 'html', 'stream-sessions.html'); + } + if (urlPath === '/async-jobs') { + return path.join(ROOT_DIR, 'html', 'async-jobs.html'); + } if (urlPath === '/settings') { return path.join(ROOT_DIR, 'html', 'settings.html'); } @@ -158,6 +164,14 @@ const server = http.createServer((request, response) => { redirect(response, '/usage'); return; } + if (urlPath === '/html/stream-sessions.html') { + redirect(response, '/stream-sessions'); + return; + } + if (urlPath === '/html/async-jobs.html') { + redirect(response, '/async-jobs'); + return; + } if (urlPath === '/html/settings.html') { redirect(response, '/settings'); return; diff --git a/apps/ui/scripts/stream-fixture-server.js b/apps/ui/scripts/stream-fixture-server.js new file mode 100644 index 0000000..7580144 --- /dev/null +++ b/apps/ui/scripts/stream-fixture-server.js @@ -0,0 +1,145 @@ +const crypto = require('crypto'); +const http = require('http'); + +const PORT = Number(process.env.CRANK_E2E_STREAM_FIXTURE_PORT || 3310); + +function writeJson(response, status, payload) { + response.writeHead(status, { + 'Content-Type': 'application/json; charset=utf-8', + 'Cache-Control': 'no-store', + }); + response.end(JSON.stringify(payload)); +} + +const server = http.createServer((request, response) => { + if (request.url === '/health') { + writeJson(response, 200, { ok: true }); + return; + } + + if (request.url === '/sse/logs') { + response.writeHead(200, { + 'Content-Type': 'text/event-stream; charset=utf-8', + 'Cache-Control': 'no-cache', + Connection: 'keep-alive', + }); + + const events = [ + { level: 'info', message: 'billing started', cursor: 'c1' }, + { level: 'warn', message: 'cache warmup slow', cursor: 'c2' }, + { level: 'error', message: 'invoice timeout', cursor: 'c3' }, + ]; + + let index = 0; + const timer = setInterval(() => { + if (index >= events.length) { + response.write('event: done\n'); + response.write('data: {"done":true}\n\n'); + clearInterval(timer); + response.end(); + return; + } + + response.write('event: message\n'); + response.write(`data: ${JSON.stringify(events[index])}\n\n`); + index += 1; + }, 150); + + request.on('close', () => { + clearInterval(timer); + }); + + return; + } + + if (request.url === '/snapshot/metrics') { + writeJson(response, 200, { + summary: { + service: 'billing', + error_rate: 0.12, + }, + items: [ + { timestamp: '2026-04-06T10:00:00Z', cpu: 0.41, memory: 0.67 }, + { timestamp: '2026-04-06T10:00:05Z', cpu: 0.39, memory: 0.65 }, + ], + done: true, + }); + return; + } + + writeJson(response, 404, { error: 'not_found' }); +}); + +function writeWebSocketFrame(socket, payload) { + const body = Buffer.from(payload, 'utf8'); + const header = []; + header.push(0x81); + if (body.length < 126) { + header.push(body.length); + } else if (body.length < 65536) { + header.push(126, (body.length >> 8) & 0xff, body.length & 0xff); + } else { + throw new Error('fixture payload is unexpectedly large'); + } + + socket.write(Buffer.concat([Buffer.from(header), body])); +} + +function acceptWebSocket(request, socket) { + const key = request.headers['sec-websocket-key']; + if (!key) { + socket.destroy(); + return; + } + + const accept = crypto + .createHash('sha1') + .update(`${key}258EAFA5-E914-47DA-95CA-C5AB0DC85B11`) + .digest('base64'); + + socket.write( + [ + 'HTTP/1.1 101 Switching Protocols', + 'Upgrade: websocket', + 'Connection: Upgrade', + `Sec-WebSocket-Accept: ${accept}`, + '\r\n', + ].join('\r\n'), + ); + + const messages = [ + { type: 'tick', seq: 1, value: 101 }, + { type: 'tick', seq: 2, value: 102 }, + { type: 'tick', seq: 3, value: 103 }, + ]; + + let index = 0; + const timer = setInterval(() => { + if (index >= messages.length) { + clearInterval(timer); + socket.end(); + return; + } + writeWebSocketFrame(socket, JSON.stringify(messages[index])); + index += 1; + }, 150); + + socket.on('close', () => clearInterval(timer)); + socket.on('end', () => clearInterval(timer)); + socket.on('error', () => clearInterval(timer)); +} + +server.on('upgrade', (request, socket, head) => { + if (request.url !== '/events') { + socket.destroy(); + return; + } + if (head && head.length) { + socket.unshift(head); + } + acceptWebSocket(request, socket); +}); + +server.listen(PORT, '127.0.0.1', () => { + console.log(`Streaming fixtures listening on http://127.0.0.1:${PORT}`); +}); diff --git a/crates/crank-adapter-grpc/examples/stream_fixture.rs b/crates/crank-adapter-grpc/examples/stream_fixture.rs new file mode 100644 index 0000000..e67bf67 --- /dev/null +++ b/crates/crank-adapter-grpc/examples/stream_fixture.rs @@ -0,0 +1,22 @@ +use crank_adapter_grpc::test_support::{EchoServiceImpl, echo}; +use tokio::net::TcpListener; +use tonic::transport::Server; + +#[tokio::main] +async fn main() -> Result<(), Box> { + let bind = std::env::var("CRANK_E2E_GRPC_FIXTURE_BIND") + .unwrap_or_else(|_| "127.0.0.1:3311".to_owned()); + let listener = TcpListener::bind(&bind).await?; + let incoming = tonic::transport::server::TcpIncoming::from(listener); + + println!("gRPC stream fixture listening on http://{bind}"); + + Server::builder() + .add_service(echo::echo_service_server::EchoServiceServer::new( + EchoServiceImpl, + )) + .serve_with_incoming(incoming) + .await?; + + Ok(()) +}