From c12c432e71272ddce9b7a6f87327af7d1ba24440 Mon Sep 17 00:00:00 2001 From: DavidVentura Date: Wed, 31 Jan 2024 16:11:16 +0100 Subject: [PATCH] refactor for multiple sessions --- discovery.ts | 50 +++++++++++++++++++++++++++++ handlers.ts | 24 ++++++++------ http_server.ts | 64 ++++++++++++++++++++++++++----------- impl.ts | 13 ++++++++ options.ts | 4 +++ session.ts | 87 ++++++++++++++++++++------------------------------ shim.ts | 5 +++ 7 files changed, 167 insertions(+), 80 deletions(-) create mode 100644 discovery.ts create mode 100644 options.ts diff --git a/discovery.ts b/discovery.ts new file mode 100644 index 0000000..394505f --- /dev/null +++ b/discovery.ts @@ -0,0 +1,50 @@ +import EventEmitter from "node:events"; +import { create_LanSearch, parse_PunchPkt } from "./impl.js"; +import { createSocket, RemoteInfo } from "node:dgram"; +import { Commands } from "./datatypes.js"; +import { hexdump } from "./hexdump.js"; +import { opt } from "./options.js"; + +const handleIncomingPunch = (msg: Buffer, ee: EventEmitter, rinfo: RemoteInfo, options: opt) => { + const ab = new Uint8Array(msg).buffer; + const dv = new DataView(ab); + const cmd_id = dv.readU16(); + if (cmd_id != Commands.PunchPkt) { + return; + } + if (options.debug) { + console.log("Discovery got a PunchPkt"); + } + ee.emit("discover", rinfo, parse_PunchPkt(dv)); +}; + +export const discoverDevices = (options: opt): EventEmitter => { + const sock = createSocket("udp4"); + //const BCAST_IP = "192.168.1.255"; + const BCAST_IP = "192.168.40.101"; + const SEND_PORT = 32108; + const ee = new EventEmitter(); + + sock.on("error", (err) => { + console.error(`sock error:\n${err.stack}`); + sock.close(); + }); + + sock.on("message", (msg, rinfo) => handleIncomingPunch(msg, ee, rinfo, options)); + + sock.on("listening", () => { + sock.setBroadcast(true); + console.log("Searching for devices.."); + + let buf = create_LanSearch(); + setInterval(() => { + console.log("."); + sock.send(new Uint8Array(buf.buffer), SEND_PORT, BCAST_IP); + }, 2000); + sock.send(new Uint8Array(buf.buffer), SEND_PORT, BCAST_IP); + }); + + sock.bind(); + + return ee; +}; diff --git a/handlers.ts b/handlers.ts index d490f5e..6620476 100644 --- a/handlers.ts +++ b/handlers.ts @@ -3,7 +3,7 @@ import { RemoteInfo } from "node:dgram"; import { Commands, CommandsByValue, ControlCommands } from "./datatypes.js"; import { XqBytesDec } from "./func_replacements.js"; import { hexdump } from "./hexdump.js"; -import { create_P2pRdy, SendListWifi, SendUsrChk } from "./impl.js"; +import { parse_PunchPkt, create_P2pRdy, SendListWifi, SendUsrChk, DevSerial } from "./impl.js"; import { Session } from "./session.js"; import { u16_swap, u32_swap } from "./utils.js"; @@ -33,15 +33,21 @@ export const handle_P2PRdy = (session: Session, _: DataView) => { session.send(b); }; +export const makePunchPkt = (dev: DevSerial): DataView => { + const len = dev.prefix.length + dev.suffix.length + 8; + const outbuf = new DataView(new Uint8Array(0x14).buffer); // 8 = serial u64 + console.log(len, dev.prefix.length, dev.serialU64, dev.suffix.length); + console.log(dev.prefix, dev.serialU64, dev.suffix); + outbuf.add(0).writeString(dev.prefix); + outbuf.add(4).writeU64(dev.serialU64); + outbuf.add(8 + dev.prefix.length).writeString(dev.suffix); + return create_P2pRdy(outbuf); +}; export const handle_PunchPkt = (session: Session, dv: DataView, rinfo: RemoteInfo) => { - const punchCmd = dv.readU16(); - const len = dv.add(2).readU16(); - const prefix = dv.add(4).readString(4); - const serial = dv.add(8).readU64().toString(); - const suffix = dv.add(16).readString(len - 16 + 4); // 16 = offset, +4 header - // f141 20 BATC 609531 EXLVS - session.eventEmitter.emit("connect", prefix.toString() + serial + suffix.toString(), rinfo); - session.send(create_P2pRdy(dv.add(4).readByteArray(len))); + const dev = parse_PunchPkt(dv); + session.eventEmitter.emit("connect", dev.devId, rinfo); + console.log("connect????"); + session.send(makePunchPkt(dev)); }; export const createResponseForControlCommand = (session: Session, dv: DataView): DataView[] => { diff --git a/http_server.ts b/http_server.ts index cf0517b..0c9bde4 100644 --- a/http_server.ts +++ b/http_server.ts @@ -1,31 +1,26 @@ import { createWriteStream } from "node:fs"; import http from "node:http"; -import { Handlers, makeSession } from "./session.js"; +import { Handlers, makeSession, Session } from "./session.js"; +import { discoverDevices } from "./discovery.js"; +import { DevSerial } from "./impl.js"; -const s = makeSession(Handlers, { debug: false, ansi: false }); -const withAudio = false; +const opts = { debug: false, ansi: false }; let BOUNDARY = "a very good boundary line"; let responses = []; - -s.eventEmitter.on("frame", (frame: Buffer) => { - let s = `--${BOUNDARY}\r\n`; - s += "Content-Type: image/jpeg\r\n\r\n"; - responses.forEach((res) => { - res.write(Buffer.from(s)); - res.write(frame); - }); -}); - -if (withAudio) { - const audioFd = createWriteStream(`audio.pcm`); - s.eventEmitter.on("audio", (frame: Buffer) => { - audioFd.write(frame); - }); -} +let sessions: Record = {}; const server = http.createServer((req, res) => { + let devId = req.url.slice(1); + console.log("requested for", devId); + let s = sessions[devId]; + + if (s === undefined) { + res.writeHead(400); + res.end("invalid ID"); + return; + } if (!s.connected) { res.writeHead(400); res.end("Nothing online"); @@ -35,4 +30,35 @@ const server = http.createServer((req, res) => { responses.push(res); }); +let devEv = discoverDevices(opts); +devEv.on("discover", (rinfo, dev: DevSerial) => { + if (dev.devId in sessions) { + console.log(`ignoring ${dev.devId} - ${rinfo.address}`); + return; + } + console.log(`discovered ${dev.devId} - ${rinfo.address}`); + const s = makeSession(Handlers, dev, rinfo, opts); + const withAudio = false; + s.eventEmitter.on("frame", (frame: Buffer) => { + let s = `--${BOUNDARY}\r\n`; + s += "Content-Type: image/jpeg\r\n\r\n"; + responses.forEach((res) => { + res.write(Buffer.from(s)); + res.write(frame); + }); + }); + + s.eventEmitter.on("disconnect", () => { + console.log("deleting from sessions"); + sessions[dev.devId] = undefined; + }); + if (withAudio) { + const audioFd = createWriteStream(`audio.pcm`); + s.eventEmitter.on("audio", (frame: Buffer) => { + audioFd.write(frame); + }); + } + sessions[dev.devId] = s; +}); + server.listen(1234); diff --git a/impl.ts b/impl.ts index 44886ea..c0043e7 100644 --- a/impl.ts +++ b/impl.ts @@ -180,3 +180,16 @@ export const create_P2pClose = (): DataView => { outbuf.add(2).writeU16(0); return outbuf; }; + +export type DevSerial = { prefix: string; serial: string; suffix: string; serialU64: bigint; devId: string }; +export const parse_PunchPkt = (dv: DataView): DevSerial => { + const punchCmd = dv.readU16(); + const len = dv.add(2).readU16(); + const prefix = dv.add(4).readString(4); + const serialU64 = dv.add(8).readU64(); + const serial = serialU64.toString(); + const suffix = dv.add(16).readString(len - 16 + 4); // 16 = offset, +4 header + const devId = prefix + serial + suffix; + + return { prefix, serial, suffix, serialU64, devId }; +}; diff --git a/options.ts b/options.ts new file mode 100644 index 0000000..66b5f86 --- /dev/null +++ b/options.ts @@ -0,0 +1,4 @@ +export type opt = { + debug: boolean; + ansi: boolean; +}; diff --git a/session.ts b/session.ts index 9e84082..18f04e0 100644 --- a/session.ts +++ b/session.ts @@ -1,14 +1,22 @@ import { createSocket, RemoteInfo } from "node:dgram"; -import { create_LanSearch, create_P2pAlive } from "./impl.js"; +import { create_P2pAlive, DevSerial } from "./impl.js"; import { Commands, CommandsByValue } from "./datatypes.js"; -import { handle_P2PAlive, handle_PunchPkt, handle_P2PRdy, handle_Drw, notImpl, noop } from "./handlers.js"; +import { + handle_P2PAlive, + handle_PunchPkt, + handle_P2PRdy, + handle_Drw, + notImpl, + noop, + makePunchPkt, +} from "./handlers.js"; import { hexdump } from "./hexdump.js"; import EventEmitter from "node:events"; import { SendVideoResolution, SendStartVideo, SendWifiDetails } from "./impl.js"; +import { opt } from "./options.js"; export type Session = { send: (msg: DataView) => void; - broadcast: (msg: DataView) => void; outgoingCommandId: number; ticket: number[]; eventEmitter: EventEmitter; @@ -21,11 +29,6 @@ export type Session = { export type PacketHandler = (session: Session, dv: DataView, rinfo: RemoteInfo) => void; -type opt = { - debug: boolean; - ansi: boolean; -}; - type msgCb = ( session: Session, handlers: Record, @@ -46,7 +49,12 @@ const handleIncoming: msgCb = (session, handlers, msg, rinfo, options) => { session.lastReceivedPacket = Date.now(); }; -export const makeSession = (handlers: Record, options: opt): Session => { +export const makeSession = ( + handlers: Record, + dev: DevSerial, + ra: RemoteInfo, + options: opt, +): Session => { const sock = createSocket("udp4"); sock.on("error", (err) => { @@ -57,33 +65,31 @@ export const makeSession = (handlers: Record handleIncoming(session, handlers, msg, rinfo, options)); sock.on("listening", () => { - const address = sock.address(); - console.log(`sock listening ${address.address}:${address.port}`); - sock.setBroadcast(true); - - // ther should be a better way of executing periodic status update - // requests per device - let buf = create_LanSearch(); - const int = setInterval(() => { - session.broadcast(buf); - }, 2000); - session.timers.push(int); - session.broadcast(buf); + const buf = makePunchPkt(dev); + session.send(buf); }); - //const BCAST_IP = "192.168.1.255"; - const BCAST_IP = "192.168.40.101"; const SEND_PORT = 32108; sock.bind(); + const sessTimer = setInterval(() => { + const delta = Date.now() - session.lastReceivedPacket; + if (delta > 600) { + let buf = create_P2pAlive(); + session.send(buf); + } + if (delta > 8000) { + session.eventEmitter.emit("disconnect"); + } + }, 400); const session: Session = { outgoingCommandId: 0, ticket: [0, 0, 0, 0], lastReceivedPacket: 0, eventEmitter: new EventEmitter(), - connected: false, - timers: [], - devName: "", + connected: true, + timers: [sessTimer], + devName: dev.serial, send: (msg: DataView) => { const raw = msg.readU16(); const cmd = CommandsByValue[raw]; @@ -95,40 +101,17 @@ export const makeSession = (handlers: Record { - sock.send(new Uint8Array(msg.buffer), SEND_PORT, BCAST_IP); - }, - dst_ip: BCAST_IP, + dst_ip: ra.address, }; - session.eventEmitter.on("disconnect", (name: string, rinfo: RemoteInfo) => { - console.log(`Disconnected from ${name} - ${rinfo.address}`); + session.eventEmitter.on("disconnect", () => { + console.log(`Disconnected from ${session.devName} - ${session.dst_ip}`); session.dst_ip = "0.0.0.0"; session.connected = false; session.timers.forEach((x) => clearInterval(x)); session.timers = []; }); - session.eventEmitter.on("connect", (name: string, rinfo: RemoteInfo) => { - console.log(`Connected to ${name} - ${rinfo.address}`); - session.outgoingCommandId = 0; - session.dst_ip = rinfo.address; - session.connected = true; - session.devName = name; - - const int = setInterval(() => { - const delta = Date.now() - session.lastReceivedPacket; - if (delta > 600) { - let buf = create_P2pAlive(); - session.send(buf); - } - if (delta > 8000) { - session.eventEmitter.emit("disconnect", name, rinfo); - } - }, 400); - session.timers.push(int); - }); - session.eventEmitter.on("login", () => { console.log(`Logged in - ${session.devName}`); startVideoStream(session); diff --git a/shim.ts b/shim.ts index 141604f..ea8fa1d 100644 --- a/shim.ts +++ b/shim.ts @@ -55,11 +55,16 @@ DataView.prototype.readString = function (len) { if (nullByte !== -1) return s.substring(0, nullByte); return s; }; +DataView.prototype.writeString = function (str) { + const bytes = [...str].map((_, i) => str.charCodeAt(i)); + return this.writeByteArray(bytes); +}; declare global { interface DataView { add(offset: number): DataView; readByteArray(len: number): DataView; + writeString(str: string): void; readString(len: number): string; readU16(): number; readU16LE(): number;