From 1f42637fcf9d70d4c0c09fd4652a0f429b7d8777 Mon Sep 17 00:00:00 2001 From: EugeneTes Date: Tue, 25 Aug 2026 10:39:22 +0000 Subject: [PATCH] Add nats.ws smoke test (connect, pub/sub, req/rep, JetStream) --- test-nats.mjs | 76 +++++++++++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 76 insertions(+) create mode 100644 test-nats.mjs diff --git a/test-nats.mjs b/test-nats.mjs new file mode 100644 index 0000000..7d608cd --- /dev/null +++ b/test-nats.mjs @@ -0,0 +1,76 @@ +// Small smoke test for the deployed NATS-over-WebSocket server. +// Verifies: connect, pub/sub roundtrip, request/reply, JetStream stream lifecycle. +// +// Setup (once): +// npm install nats.ws ws +// +// Usage: +// NATS_URL=wss://nats.tes.gd NATS_USER=admin NATS_PASSWORD=... node test-nats.mjs + +// Polyfill the browser WebSocket global; nats.ws requires it. +if (typeof globalThis.WebSocket === 'undefined') { + const { default: WS } = await import('ws'); + globalThis.WebSocket = WS; +} +import { connect, StringCodec } from 'nats.ws'; + +const url = process.env.NATS_URL ?? 'wss://nats.tes.gd'; +const user = process.env.NATS_USER ?? 'admin'; +const pass = process.env.NATS_PASSWORD; + +if (!pass) { console.error('NATS_PASSWORD is required'); process.exit(1); } + +const sc = StringCodec(); +const fail = (msg, err) => { console.error(`FAIL: ${msg}`, err ?? ''); process.exit(1); }; + +console.log(`→ connecting to ${url} as ${user}`); +let nc; +try { + nc = await connect({ servers: url, user, pass, timeout: 5000 }); +} catch (e) { fail('connect', e); } +console.log(`✓ connected server=${nc.info.server_name} version=${nc.info.version} jetstream=${nc.info.jetstream}`); + +// 1. Pub/sub roundtrip +const subj = `test.${Date.now()}`; +const sub = nc.subscribe(subj, { max: 1 }); +const gotPub = (async () => { + for await (const m of sub) return sc.decode(m.data); +})(); +nc.publish(subj, sc.encode('hello nats')); +const payload = await Promise.race([gotPub, new Promise((_, r) => setTimeout(() => r(new Error('timeout')), 3000))]); +if (payload !== 'hello nats') fail(`pub/sub payload mismatch: "${payload}"`); +console.log(`✓ pub/sub subject=${subj} payload="${payload}"`); + +// 2. Request/reply +const svcSub = nc.subscribe('svc.echo', { max: 1 }); +(async () => { + for await (const m of svcSub) m.respond(sc.encode(`echo:${sc.decode(m.data)}`)); +})(); +const reply = await nc.request('svc.echo', sc.encode('ping'), { timeout: 3000 }); +if (sc.decode(reply.data) !== 'echo:ping') fail(`req/rep payload mismatch: "${sc.decode(reply.data)}"`); +console.log(`✓ request/rep reply="${sc.decode(reply.data)}"`); + +// 3. JetStream: create stream, publish, consume, delete +const jsm = await nc.jetstreamManager(); +const streamName = `TEST_${Date.now()}`; +const streamSubj = `js.${streamName}.>`; +await jsm.streams.add({ name: streamName, subjects: [streamSubj], storage: 'file' }); +const js = nc.jetstream(); +const ack = await js.publish(`${streamSubj.replace('>', 'hello')}`, sc.encode('jetstream-payload')); +if (ack.seq !== 1) fail(`js publish seq expected 1, got ${ack.seq}`); +console.log(`✓ jetstream stream=${streamName} seq=${ack.seq} domain=${ack.domain ?? '-'}`); + +const consumerName = `TEST_C_${Date.now()}`; +await jsm.consumers.add(streamName, { durable_name: consumerName, ack_policy: 'explicit' }); +const consumer = await js.consumers.get(streamName, consumerName); +const iter = await consumer.fetch({ max_messages: 1, expires: 2000 }); +let jsPayload; +for await (const m of iter) { jsPayload = sc.decode(m.data); m.ack(); } +if (jsPayload !== 'jetstream-payload') fail(`jetstream payload mismatch: "${jsPayload}"`); +console.log(`✓ jetstream consumed payload="${jsPayload}"`); + +await jsm.streams.delete(streamName); +console.log(`✓ cleanup stream=${streamName} deleted`); + +await nc.drain(); +console.log('✓ ALL CHECKS PASSED');