commit a2f5814f1f414f501db9b4cf8d300f54f6932664 Author: Peter Pötzi Date: Fri Apr 10 21:34:16 2026 +0200 Initial commit diff --git a/README.md b/README.md new file mode 100644 index 0000000..fabb204 --- /dev/null +++ b/README.md @@ -0,0 +1,34 @@ +# cits-wireshark-bridge + +Node.js script that bridges raw 802.11 packets on NATS into `tshark`, and publishes the decoded JSON back to NATS. + +We use it for our opentrafficmap.org project. + +## Requirements + +- Node.js 18 or newer +- `nats` for Node.js (`npm install`) +- `tshark` installed and available in `PATH` +- A reachable NATS server + +## Environment + +- `NATS_URL`: NATS server URL, for example `nats://127.0.0.1:4222` +- `NATS_USERNAME`, `NATS_PASSWORD`: optional credentials +- `INPUT_SUBJECT`: subscription subject, default `its.*.packet` +- `OUTPUT_SUBJECT`: fixed publish subject; if unset, the bridge publishes to `.json` +- `OUTPUT_SUFFIX`: suffix used when `OUTPUT_SUBJECT` is unset, default `.json` +- `PCAP_SNAPLEN`: PCAP snapshot length, default `65535` +- `PCAP_LINKTYPE`: PCAP link type, default `105` +- `TSHARK_BIN`: path to `tshark`, default `tshark` +- `TSHARK_PROTOCOL_FILTER`: optional `-J` protocol filter for `tshark` +- `TSHARK_INCLUDE_HEX`: set to `1`, `true`, or `yes` to include hex output +- `VERBOSE`: set to `1`, `true`, or `yes` for extra logging + +## Run + +```bash +NATS_URL=nats://127.0.0.1:4222 \ +INPUT_SUBJECT='its.*.packet' \ +node its-bridge.js +``` diff --git a/its-bridge.js b/its-bridge.js new file mode 100644 index 0000000..295e4bd --- /dev/null +++ b/its-bridge.js @@ -0,0 +1,353 @@ +#!/usr/bin/env node +'use strict'; + +const readline = require('readline'); +const { spawn } = require('child_process'); + +const NATS_URL = process.env.NATS_URL || 'nats://127.0.0.1:4222'; +const NATS_USERNAME = process.env.NATS_USERNAME || ''; +const NATS_PASSWORD = process.env.NATS_PASSWORD || ''; +const INPUT_SUBJECT = process.env.INPUT_SUBJECT || 'its.*.packet'; +const OUTPUT_SUBJECT = process.env.OUTPUT_SUBJECT || ''; +const OUTPUT_SUFFIX = process.env.OUTPUT_SUFFIX || '.json'; +const PCAP_SNAPLEN = parseIntegerEnv('PCAP_SNAPLEN', 65535, { min: 1 }); +const PCAP_LINKTYPE = parseIntegerEnv('PCAP_LINKTYPE', 105, { min: 0 }); // LINKTYPE_IEEE802_11 +const TSHARK_BIN = process.env.TSHARK_BIN || 'tshark'; +const TSHARK_PROTOCOL_FILTER = process.env.TSHARK_PROTOCOL_FILTER || ''; +const INCLUDE_HEX = parseBooleanEnv('TSHARK_INCLUDE_HEX'); +const VERBOSE = parseBooleanEnv('VERBOSE'); + +let nc; +let sub; +let jc; +let tshark; +let rl; +let shuttingDown = false; +let pcapHeaderWritten = false; +let writePaused = false; +let pcapWriteQueue = []; +let pendingPackets = []; +let seq = 0; +let stats = { + in: 0, + pcapWritten: 0, + jsonPublished: 0, + tsharkMetaLines: 0, + parseErrors: 0, + droppedBecauseNoMeta: 0, +}; + +function parseBooleanEnv(name, defaultValue = false) { + const raw = process.env[name]; + if (raw == null || raw === '') return defaultValue; + + if (/^(1|true|yes)$/i.test(raw)) return true; + if (/^(0|false|no)$/i.test(raw)) return false; + + throw new Error(`${name} must be one of: 1, 0, true, false, yes, no`); +} + +function parseIntegerEnv(name, defaultValue, { min = Number.MIN_SAFE_INTEGER, max = Number.MAX_SAFE_INTEGER } = {}) { + const raw = process.env[name]; + if (raw == null || raw === '') return defaultValue; + + if (!/^-?\d+$/.test(raw)) { + throw new Error(`${name} must be an integer`); + } + + const value = Number(raw); + if (!Number.isSafeInteger(value) || value < min || value > max) { + throw new Error(`${name} must be an integer between ${min} and ${max}`); + } + + return value; +} + +function log(...args) { + process.stdout.write(`${new Date().toISOString()} ${args.join(' ')}\n`); +} + +function logErr(...args) { + process.stderr.write(`${new Date().toISOString()} ${args.join(' ')}\n`); +} + +function normalizePayload(data) { + if (Buffer.isBuffer(data)) return data; + if (data instanceof Uint8Array) return Buffer.from(data); + throw new TypeError(`Unsupported payload type: ${typeof data}`); +} + +function makePcapGlobalHeader() { + const hdr = Buffer.alloc(24); + hdr.writeUInt32LE(0xa1b2c3d4, 0); + hdr.writeUInt16LE(2, 4); + hdr.writeUInt16LE(4, 6); + hdr.writeInt32LE(0, 8); + hdr.writeUInt32LE(0, 12); + hdr.writeUInt32LE(PCAP_SNAPLEN, 16); + hdr.writeUInt32LE(PCAP_LINKTYPE, 20); + return hdr; +} + +function makePcapRecord(packetBuf, tsMs) { + const tsSec = Math.floor(tsMs / 1000); + const tsUsec = (tsMs % 1000) * 1000; + const inclLen = Math.min(packetBuf.length, PCAP_SNAPLEN); + const pkt = inclLen === packetBuf.length ? packetBuf : packetBuf.subarray(0, inclLen); + + const recHdr = Buffer.alloc(16); + recHdr.writeUInt32LE(tsSec, 0); + recHdr.writeUInt32LE(tsUsec, 4); + recHdr.writeUInt32LE(inclLen, 8); + recHdr.writeUInt32LE(packetBuf.length, 12); + + return { recHdr, pkt }; +} + +function ensurePcapHeaderQueued() { + if (pcapHeaderWritten) return; + pcapHeaderWritten = true; + pcapWriteQueue.push(makePcapGlobalHeader()); +} + +function pumpPcapToTshark() { + if (!tshark || !tshark.stdin || writePaused) return; + + while (pcapWriteQueue.length > 0) { + const chunk = pcapWriteQueue[0]; + const ok = tshark.stdin.write(chunk); + if (!ok) { + writePaused = true; + return; + } + pcapWriteQueue.shift(); + } +} + +function queuePacketForTshark(sourceSubject, packetBuf) { + ensurePcapHeaderQueued(); + + const now = Date.now(); + const packetSeq = ++seq; + const { recHdr, pkt } = makePcapRecord(packetBuf, now); + + pendingPackets.push({ + seq: packetSeq, + source_subject: sourceSubject, + received_ts: new Date(now).toISOString(), + raw_len: packetBuf.length, + incl_len: pkt.length, + }); + + pcapWriteQueue.push(recHdr, pkt); + stats.pcapWritten += 1; + pumpPcapToTshark(); +} + +function isEkMetaLine(obj) { + if (!obj || typeof obj !== 'object' || Array.isArray(obj)) return false; + return Boolean(obj.index || obj.create || obj.update || obj.delete); +} + +function outputSubjectFor(sourceSubject) { + if (OUTPUT_SUBJECT) return OUTPUT_SUBJECT; + return `${sourceSubject}${OUTPUT_SUFFIX}`; +} + +async function publishDecodedPacket(decodedObj) { + const meta = pendingPackets.shift(); + if (!meta) { + stats.droppedBecauseNoMeta += 1; + if (VERBOSE) logErr('decoded packet without pending metadata'); + return; + } + + const envelope = { + seq: meta.seq, + source_subject: meta.source_subject, + received_ts: meta.received_ts, + published_ts: new Date().toISOString(), + raw_len: meta.raw_len, + incl_len: meta.incl_len, + decoded: decodedObj, + }; + + const outSubject = outputSubjectFor(meta.source_subject); + nc.publish(outSubject, jc.encode(envelope)); + stats.jsonPublished += 1; +} + +async function handleTsharkLine(line) { + const trimmed = line.trim(); + if (!trimmed) return; + + let obj; + try { + obj = JSON.parse(trimmed); + } catch (err) { + stats.parseErrors += 1; + logErr('failed to parse tshark json line:', err.message); + if (VERBOSE) { + logErr(trimmed); + } + return; + } + + if (isEkMetaLine(obj)) { + stats.tsharkMetaLines += 1; + return; + } + + await publishDecodedPacket(obj); +} + +function buildTsharkArgs() { + const args = ['-i', '-', '-l', '-n', '-T', 'ek']; + if (TSHARK_PROTOCOL_FILTER) { + args.push('-J', TSHARK_PROTOCOL_FILTER); + } + if (INCLUDE_HEX) { + args.push('-x'); + } + return args; +} + +function attachTsharkHandlers() { + tshark.on('error', (err) => { + if (shuttingDown) return; + logErr(`failed to start tshark (${TSHARK_BIN}):`, err.message || String(err)); + process.exit(1); + }); + + tshark.stdin.on('drain', () => { + writePaused = false; + pumpPcapToTshark(); + }); + + tshark.stdin.on('error', (err) => { + if (shuttingDown && err && err.code === 'EPIPE') return; + logErr('tshark stdin error:', err.message || String(err)); + }); + + tshark.stderr.on('data', (chunk) => { + const text = String(chunk); + if (!text.trim()) return; + logErr(`[tshark] ${text.trimEnd()}`); + }); + + tshark.on('exit', (code, signal) => { + if (shuttingDown) return; + logErr(`tshark exited unexpectedly code=${code} signal=${signal}`); + process.exit(code || 1); + }); + + rl = readline.createInterface({ input: tshark.stdout, crlfDelay: Infinity }); + rl.on('line', (line) => { + handleTsharkLine(line).catch((err) => { + logErr('error handling tshark line:', err.stack || String(err)); + }); + }); +} + +async function shutdown(reason) { + if (shuttingDown) return; + shuttingDown = true; + log(`shutdown: ${reason}`); + + try { + if (sub) sub.unsubscribe(); + } catch (err) { + logErr('unsubscribe failed:', err.message || String(err)); + } + + try { + if (tshark && tshark.stdin && !tshark.stdin.destroyed) { + tshark.stdin.end(); + } + } catch (err) { + logErr('closing tshark stdin failed:', err.message || String(err)); + } + + try { + if (nc) { + await nc.flush(); + await nc.drain(); + } + } catch (err) { + logErr('nats drain failed:', err.message || String(err)); + } + + log(`stats in=${stats.in} pcapWritten=${stats.pcapWritten} jsonPublished=${stats.jsonPublished} metaLines=${stats.tsharkMetaLines} parseErrors=${stats.parseErrors} pending=${pendingPackets.length}`); + process.exit(0); +} + +async function main() { + const { connect, JSONCodec } = await import('nats'); + jc = JSONCodec(); + + const tsharkArgs = buildTsharkArgs(); + log(`starting tshark: ${TSHARK_BIN} ${tsharkArgs.join(' ')}`); + log(`nats input=${INPUT_SUBJECT} output=${OUTPUT_SUBJECT || `${OUTPUT_SUFFIX}`}`); + log(`pcap linktype=${PCAP_LINKTYPE} snaplen=${PCAP_SNAPLEN}`); + + tshark = spawn(TSHARK_BIN, tsharkArgs, { + stdio: ['pipe', 'pipe', 'pipe'], + }); + attachTsharkHandlers(); + + nc = await connect({ + servers: NATS_URL, + user: NATS_USERNAME, + pass: NATS_PASSWORD, + }); + log(`connected to ${nc.getServer()}`); + + sub = nc.subscribe(INPUT_SUBJECT, { + callback: (err, msg) => { + if (err) { + logErr('subscription error:', err.message || String(err)); + return; + } + + try { + const packetBuf = normalizePayload(msg.data); + stats.in += 1; + queuePacketForTshark(msg.subject, packetBuf); + } catch (innerErr) { + logErr('failed to queue packet:', innerErr.stack || String(innerErr)); + } + }, + }); + + const closed = nc.closed(); + closed.then((err) => { + if (err) { + logErr('nats connection closed with error:', err.message || String(err)); + } else { + log('nats connection closed'); + } + shutdown('nats closed').catch((shutdownErr) => { + logErr('shutdown after nats close failed:', shutdownErr.stack || String(shutdownErr)); + process.exit(1); + }); + }); +} + +process.on('SIGINT', () => { + shutdown('SIGINT').catch((err) => { + logErr('shutdown failed:', err.stack || String(err)); + process.exit(1); + }); +}); + +process.on('SIGTERM', () => { + shutdown('SIGTERM').catch((err) => { + logErr('shutdown failed:', err.stack || String(err)); + process.exit(1); + }); +}); + +main().catch((err) => { + logErr(err.stack || String(err)); + process.exit(1); +}); diff --git a/package.json b/package.json new file mode 100644 index 0000000..212f054 --- /dev/null +++ b/package.json @@ -0,0 +1,22 @@ +{ + "name": "cits-wireshark-bridge", + "version": "0.1.0", + "description": "Pipes raw 802.11 packets from a NATS connection into wireshark (tshark) and sends a json version of the packet back to the NATS server", + "main": "its-bridge.js", + "bin": { + "cits-wireshark-bridge": "./its-bridge.js" + }, + "files": [ + "its-bridge.js", + "README.md" + ], + "scripts": { + "start": "node its-bridge.js" + }, + "engines": { + "node": ">=18" + }, + "dependencies": { + "nats": "^2.0.0" + } +}