refactor for multiple sessions
This commit is contained in:
@@ -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;
|
||||||
|
};
|
||||||
+15
-9
@@ -3,7 +3,7 @@ import { RemoteInfo } from "node:dgram";
|
|||||||
import { Commands, CommandsByValue, ControlCommands } from "./datatypes.js";
|
import { Commands, CommandsByValue, ControlCommands } from "./datatypes.js";
|
||||||
import { XqBytesDec } from "./func_replacements.js";
|
import { XqBytesDec } from "./func_replacements.js";
|
||||||
import { hexdump } from "./hexdump.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 { Session } from "./session.js";
|
||||||
import { u16_swap, u32_swap } from "./utils.js";
|
import { u16_swap, u32_swap } from "./utils.js";
|
||||||
|
|
||||||
@@ -33,15 +33,21 @@ export const handle_P2PRdy = (session: Session, _: DataView) => {
|
|||||||
session.send(b);
|
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) => {
|
export const handle_PunchPkt = (session: Session, dv: DataView, rinfo: RemoteInfo) => {
|
||||||
const punchCmd = dv.readU16();
|
const dev = parse_PunchPkt(dv);
|
||||||
const len = dv.add(2).readU16();
|
session.eventEmitter.emit("connect", dev.devId, rinfo);
|
||||||
const prefix = dv.add(4).readString(4);
|
console.log("connect????");
|
||||||
const serial = dv.add(8).readU64().toString();
|
session.send(makePunchPkt(dev));
|
||||||
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)));
|
|
||||||
};
|
};
|
||||||
|
|
||||||
export const createResponseForControlCommand = (session: Session, dv: DataView): DataView[] => {
|
export const createResponseForControlCommand = (session: Session, dv: DataView): DataView[] => {
|
||||||
|
|||||||
+45
-19
@@ -1,31 +1,26 @@
|
|||||||
import { createWriteStream } from "node:fs";
|
import { createWriteStream } from "node:fs";
|
||||||
import http from "node:http";
|
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 opts = { debug: false, ansi: false };
|
||||||
const withAudio = false;
|
|
||||||
|
|
||||||
let BOUNDARY = "a very good boundary line";
|
let BOUNDARY = "a very good boundary line";
|
||||||
let responses = [];
|
let responses = [];
|
||||||
|
let sessions: Record<string, Session> = {};
|
||||||
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);
|
|
||||||
});
|
|
||||||
}
|
|
||||||
|
|
||||||
const server = http.createServer((req, res) => {
|
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) {
|
if (!s.connected) {
|
||||||
res.writeHead(400);
|
res.writeHead(400);
|
||||||
res.end("Nothing online");
|
res.end("Nothing online");
|
||||||
@@ -35,4 +30,35 @@ const server = http.createServer((req, res) => {
|
|||||||
responses.push(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);
|
server.listen(1234);
|
||||||
|
|||||||
@@ -180,3 +180,16 @@ export const create_P2pClose = (): DataView => {
|
|||||||
outbuf.add(2).writeU16(0);
|
outbuf.add(2).writeU16(0);
|
||||||
return outbuf;
|
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 };
|
||||||
|
};
|
||||||
|
|||||||
@@ -0,0 +1,4 @@
|
|||||||
|
export type opt = {
|
||||||
|
debug: boolean;
|
||||||
|
ansi: boolean;
|
||||||
|
};
|
||||||
+35
-52
@@ -1,14 +1,22 @@
|
|||||||
import { createSocket, RemoteInfo } from "node:dgram";
|
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 { 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 { hexdump } from "./hexdump.js";
|
||||||
import EventEmitter from "node:events";
|
import EventEmitter from "node:events";
|
||||||
import { SendVideoResolution, SendStartVideo, SendWifiDetails } from "./impl.js";
|
import { SendVideoResolution, SendStartVideo, SendWifiDetails } from "./impl.js";
|
||||||
|
import { opt } from "./options.js";
|
||||||
|
|
||||||
export type Session = {
|
export type Session = {
|
||||||
send: (msg: DataView) => void;
|
send: (msg: DataView) => void;
|
||||||
broadcast: (msg: DataView) => void;
|
|
||||||
outgoingCommandId: number;
|
outgoingCommandId: number;
|
||||||
ticket: number[];
|
ticket: number[];
|
||||||
eventEmitter: EventEmitter;
|
eventEmitter: EventEmitter;
|
||||||
@@ -21,11 +29,6 @@ export type Session = {
|
|||||||
|
|
||||||
export type PacketHandler = (session: Session, dv: DataView, rinfo: RemoteInfo) => void;
|
export type PacketHandler = (session: Session, dv: DataView, rinfo: RemoteInfo) => void;
|
||||||
|
|
||||||
type opt = {
|
|
||||||
debug: boolean;
|
|
||||||
ansi: boolean;
|
|
||||||
};
|
|
||||||
|
|
||||||
type msgCb = (
|
type msgCb = (
|
||||||
session: Session,
|
session: Session,
|
||||||
handlers: Record<keyof typeof Commands, PacketHandler>,
|
handlers: Record<keyof typeof Commands, PacketHandler>,
|
||||||
@@ -46,7 +49,12 @@ const handleIncoming: msgCb = (session, handlers, msg, rinfo, options) => {
|
|||||||
session.lastReceivedPacket = Date.now();
|
session.lastReceivedPacket = Date.now();
|
||||||
};
|
};
|
||||||
|
|
||||||
export const makeSession = (handlers: Record<keyof typeof Commands, PacketHandler>, options: opt): Session => {
|
export const makeSession = (
|
||||||
|
handlers: Record<keyof typeof Commands, PacketHandler>,
|
||||||
|
dev: DevSerial,
|
||||||
|
ra: RemoteInfo,
|
||||||
|
options: opt,
|
||||||
|
): Session => {
|
||||||
const sock = createSocket("udp4");
|
const sock = createSocket("udp4");
|
||||||
|
|
||||||
sock.on("error", (err) => {
|
sock.on("error", (err) => {
|
||||||
@@ -57,33 +65,31 @@ export const makeSession = (handlers: Record<keyof typeof Commands, PacketHandle
|
|||||||
sock.on("message", (msg, rinfo) => handleIncoming(session, handlers, msg, rinfo, options));
|
sock.on("message", (msg, rinfo) => handleIncoming(session, handlers, msg, rinfo, options));
|
||||||
|
|
||||||
sock.on("listening", () => {
|
sock.on("listening", () => {
|
||||||
const address = sock.address();
|
const buf = makePunchPkt(dev);
|
||||||
console.log(`sock listening ${address.address}:${address.port}`);
|
session.send(buf);
|
||||||
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 BCAST_IP = "192.168.1.255";
|
|
||||||
const BCAST_IP = "192.168.40.101";
|
|
||||||
const SEND_PORT = 32108;
|
const SEND_PORT = 32108;
|
||||||
sock.bind();
|
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 = {
|
const session: Session = {
|
||||||
outgoingCommandId: 0,
|
outgoingCommandId: 0,
|
||||||
ticket: [0, 0, 0, 0],
|
ticket: [0, 0, 0, 0],
|
||||||
lastReceivedPacket: 0,
|
lastReceivedPacket: 0,
|
||||||
eventEmitter: new EventEmitter(),
|
eventEmitter: new EventEmitter(),
|
||||||
connected: false,
|
connected: true,
|
||||||
timers: [],
|
timers: [sessTimer],
|
||||||
devName: "",
|
devName: dev.serial,
|
||||||
send: (msg: DataView) => {
|
send: (msg: DataView) => {
|
||||||
const raw = msg.readU16();
|
const raw = msg.readU16();
|
||||||
const cmd = CommandsByValue[raw];
|
const cmd = CommandsByValue[raw];
|
||||||
@@ -95,40 +101,17 @@ export const makeSession = (handlers: Record<keyof typeof Commands, PacketHandle
|
|||||||
}
|
}
|
||||||
sock.send(new Uint8Array(msg.buffer), SEND_PORT, session.dst_ip);
|
sock.send(new Uint8Array(msg.buffer), SEND_PORT, session.dst_ip);
|
||||||
},
|
},
|
||||||
broadcast: (msg: DataView) => {
|
dst_ip: ra.address,
|
||||||
sock.send(new Uint8Array(msg.buffer), SEND_PORT, BCAST_IP);
|
|
||||||
},
|
|
||||||
dst_ip: BCAST_IP,
|
|
||||||
};
|
};
|
||||||
|
|
||||||
session.eventEmitter.on("disconnect", (name: string, rinfo: RemoteInfo) => {
|
session.eventEmitter.on("disconnect", () => {
|
||||||
console.log(`Disconnected from ${name} - ${rinfo.address}`);
|
console.log(`Disconnected from ${session.devName} - ${session.dst_ip}`);
|
||||||
session.dst_ip = "0.0.0.0";
|
session.dst_ip = "0.0.0.0";
|
||||||
session.connected = false;
|
session.connected = false;
|
||||||
session.timers.forEach((x) => clearInterval(x));
|
session.timers.forEach((x) => clearInterval(x));
|
||||||
session.timers = [];
|
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", () => {
|
session.eventEmitter.on("login", () => {
|
||||||
console.log(`Logged in - ${session.devName}`);
|
console.log(`Logged in - ${session.devName}`);
|
||||||
startVideoStream(session);
|
startVideoStream(session);
|
||||||
|
|||||||
@@ -55,11 +55,16 @@ DataView.prototype.readString = function (len) {
|
|||||||
if (nullByte !== -1) return s.substring(0, nullByte);
|
if (nullByte !== -1) return s.substring(0, nullByte);
|
||||||
return s;
|
return s;
|
||||||
};
|
};
|
||||||
|
DataView.prototype.writeString = function (str) {
|
||||||
|
const bytes = [...str].map((_, i) => str.charCodeAt(i));
|
||||||
|
return this.writeByteArray(bytes);
|
||||||
|
};
|
||||||
|
|
||||||
declare global {
|
declare global {
|
||||||
interface DataView {
|
interface DataView {
|
||||||
add(offset: number): DataView;
|
add(offset: number): DataView;
|
||||||
readByteArray(len: number): DataView;
|
readByteArray(len: number): DataView;
|
||||||
|
writeString(str: string): void;
|
||||||
readString(len: number): string;
|
readString(len: number): string;
|
||||||
readU16(): number;
|
readU16(): number;
|
||||||
readU16LE(): number;
|
readU16LE(): number;
|
||||||
|
|||||||
Reference in New Issue
Block a user