feat: scaffold streaming e2e fixtures
This commit is contained in:
@@ -4,7 +4,7 @@
|
|||||||
|
|
||||||
### `feat/streaming-e2e`
|
### `feat/streaming-e2e`
|
||||||
|
|
||||||
Status: pending
|
Status: in progress
|
||||||
|
|
||||||
DoD:
|
DoD:
|
||||||
- local fixture stack includes:
|
- local fixture stack includes:
|
||||||
|
|||||||
@@ -9,6 +9,8 @@ POSTGRES_PORT="${CRANK_E2E_POSTGRES_PORT:-55433}"
|
|||||||
ADMIN_PORT="${CRANK_E2E_ADMIN_PORT:-3301}"
|
ADMIN_PORT="${CRANK_E2E_ADMIN_PORT:-3301}"
|
||||||
MCP_PORT="${CRANK_E2E_MCP_PORT:-3302}"
|
MCP_PORT="${CRANK_E2E_MCP_PORT:-3302}"
|
||||||
UI_PORT="${CRANK_E2E_UI_PORT:-3300}"
|
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_DB="${CRANK_E2E_POSTGRES_DB:-crank}"
|
||||||
POSTGRES_USER="${CRANK_E2E_POSTGRES_USER:-crank}"
|
POSTGRES_USER="${CRANK_E2E_POSTGRES_USER:-crank}"
|
||||||
POSTGRES_PASSWORD="${CRANK_E2E_POSTGRES_PASSWORD:-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
|
kill "$(cat "$TMP_DIR/ui-server.pid")" >/dev/null 2>&1 || true
|
||||||
rm -f "$TMP_DIR/ui-server.pid"
|
rm -f "$TMP_DIR/ui-server.pid"
|
||||||
fi
|
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
|
docker rm -f "$POSTGRES_CONTAINER" >/dev/null 2>&1 || true
|
||||||
kill_port_processes "$UI_PORT"
|
kill_port_processes "$UI_PORT"
|
||||||
kill_port_processes "$ADMIN_PORT"
|
kill_port_processes "$ADMIN_PORT"
|
||||||
kill_port_processes "$MCP_PORT"
|
kill_port_processes "$MCP_PORT"
|
||||||
|
kill_port_processes "$STREAM_FIXTURE_PORT"
|
||||||
|
kill_port_processes "${GRPC_FIXTURE_BIND##*:}"
|
||||||
}
|
}
|
||||||
|
|
||||||
trap cleanup EXIT INT TERM
|
trap cleanup EXIT INT TERM
|
||||||
|
|
||||||
cleanup
|
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 \
|
docker run -d --rm \
|
||||||
--name "$POSTGRES_CONTAINER" \
|
--name "$POSTGRES_CONTAINER" \
|
||||||
-e POSTGRES_DB="$POSTGRES_DB" \
|
-e POSTGRES_DB="$POSTGRES_DB" \
|
||||||
@@ -76,6 +107,8 @@ export CRANK_BOOTSTRAP_ADMIN_DISPLAY_NAME="Crank E2E"
|
|||||||
export CRANK_DEMO_SEED="true"
|
export CRANK_DEMO_SEED="true"
|
||||||
export CRANK_PUBLIC_BASE_URL="http://127.0.0.1:$UI_PORT"
|
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_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"
|
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
|
sleep 1
|
||||||
done
|
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"
|
cd "$ROOT_DIR/apps/ui"
|
||||||
node scripts/playwright-ui-server.js >"$LOG_DIR/ui-server.log" 2>&1
|
node scripts/playwright-ui-server.js >"$LOG_DIR/ui-server.log" 2>&1
|
||||||
|
|||||||
@@ -50,6 +50,12 @@ function mapRoute(urlPath) {
|
|||||||
if (urlPath === '/usage') {
|
if (urlPath === '/usage') {
|
||||||
return path.join(ROOT_DIR, 'html', 'usage.html');
|
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') {
|
if (urlPath === '/settings') {
|
||||||
return path.join(ROOT_DIR, 'html', 'settings.html');
|
return path.join(ROOT_DIR, 'html', 'settings.html');
|
||||||
}
|
}
|
||||||
@@ -158,6 +164,14 @@ const server = http.createServer((request, response) => {
|
|||||||
redirect(response, '/usage');
|
redirect(response, '/usage');
|
||||||
return;
|
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') {
|
if (urlPath === '/html/settings.html') {
|
||||||
redirect(response, '/settings');
|
redirect(response, '/settings');
|
||||||
return;
|
return;
|
||||||
|
|||||||
@@ -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}`);
|
||||||
|
});
|
||||||
@@ -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<dyn std::error::Error>> {
|
||||||
|
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(())
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user