From 7d2160f07f6de960b4827cf9213ec5de02cca063 Mon Sep 17 00:00:00 2001 From: DavidVentura Date: Wed, 31 Jan 2024 18:06:52 +0100 Subject: [PATCH] support multiple streams concurrently --- handlers.ts | 29 +++++++++++++---------------- http_server.ts | 35 ++++++++++++++++++++++------------- session.ts | 18 ++++++++++++------ 3 files changed, 47 insertions(+), 35 deletions(-) diff --git a/handlers.ts b/handlers.ts index cb5b1d8..14efa7e 100644 --- a/handlers.ts +++ b/handlers.ts @@ -5,8 +5,6 @@ import { create_P2pRdy, SendListWifi, SendUsrChk, DevSerial } from "./impl.js"; import { Session } from "./session.js"; import { u16_swap, u32_swap } from "./utils.js"; -let curImage: Buffer | null = null; - export const notImpl = (_: Session, dv: DataView) => { const raw = dv.readU16(); const cmd = CommandsByValue[raw]; @@ -124,8 +122,6 @@ export const createResponseForControlCommand = (session: Session, dv: DataView): return []; }; -let seq = 0; -let frame_is_bad = false; const deal_with_data = (session: Session, dv: DataView) => { const pkt_len = dv.add(2).readU16(); // data @@ -152,30 +148,31 @@ const deal_with_data = (session: Session, dv: DataView) => { } else { const data = dv.add(8).readByteArray(pkt_len - 4); if (is_new_image) { - if (curImage != null && !frame_is_bad) { - session.eventEmitter.emit("frame", curImage); + if (session.curImage != null && !session.frame_is_bad) { + session.eventEmitter.emit("frame"); } - frame_is_bad = false; - curImage = Buffer.from(data.buffer); - seq = pkt_id; + session.frame_is_bad = false; + session.curImage = Buffer.from(data.buffer); + session.rcvSeqId = pkt_id; } else { - if (pkt_id <= seq) { + if (pkt_id <= session.rcvSeqId) { // retransmit return; } - if (frame_is_bad) { + if (session.frame_is_bad) { return; } - if (pkt_id > seq + 1) { + if (pkt_id > session.rcvSeqId + 1) { // missed some packets -- filling with zeroes still produces a broken // image - frame_is_bad = true; + session.frame_is_bad = true; + console.log(`Missed some packets on ${session.devName} -- skipping a frame`); return; } - seq = pkt_id; - if (curImage != null) { - curImage = Buffer.concat([curImage, Buffer.from(data.buffer)]); + session.rcvSeqId = pkt_id; + if (session.curImage != null) { + session.curImage = Buffer.concat([session.curImage, Buffer.from(data.buffer)]); } } } diff --git a/http_server.ts b/http_server.ts index 916af9d..8a3da53 100644 --- a/http_server.ts +++ b/http_server.ts @@ -1,19 +1,21 @@ +import { RemoteInfo } from "dgram"; import { createWriteStream } from "node:fs"; import http from "node:http"; -import { RemoteInfo } from "dgram"; -import { Handlers, makeSession, Session, startVideoStream } from "./session.js"; import { discoverDevices } from "./discovery.js"; import { DevSerial } from "./impl.js"; +import { Handlers, makeSession, Session, startVideoStream } from "./session.js"; +import { ServerResponse } from "http"; const opts = { debug: false, ansi: false, - discovery_ip: "192.168.40.101", //, "192.168.1.255"; + discovery_ip: "192.168.40.255", //, "192.168.1.255" + // discovery_ip: "192.168.40.101", }; let BOUNDARY = "a very good boundary line"; -let responses = []; +let responses: Record = {}; let sessions: Record = {}; const server = http.createServer((req, res) => { @@ -33,10 +35,14 @@ const server = http.createServer((req, res) => { return; } res.setHeader("Content-Type", `multipart/x-mixed-replace; boundary="${BOUNDARY}"`); - responses.push(res); + responses[devId].push(res); + res.on("close", () => { + responses[devId] = responses[devId].filter((r) => r !== res); + console.log("Conn closed, kicked"); + }); } else { res.write(``); - Object.keys(sessions).forEach((id) => res.write(`${id}`)); + Object.keys(sessions).forEach((id) => res.write(`
`)); res.write(``); res.end(); } @@ -48,21 +54,24 @@ devEv.on("discover", (rinfo: RemoteInfo, dev: DevSerial) => { console.log(`ignoring ${dev.devId} - ${rinfo.address}`); return; } + console.log(`discovered ${dev.devId} - ${rinfo.address}`); + responses[dev.devId] = []; const s = makeSession(Handlers, dev, rinfo, startVideoStream, 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); + const header = Buffer.from(`--${BOUNDARY}\r\nContent-Type: image/jpeg\r\n\r\n`); + + s.eventEmitter.on("frame", () => { + responses[dev.devId].forEach((res) => { + res.write(header); + res.write(s.curImage); }); }); s.eventEmitter.on("disconnect", () => { console.log("deleting from sessions"); - sessions[dev.devId] = undefined; + delete sessions[dev.devId]; }); if (withAudio) { const audioFd = createWriteStream(`audio.pcm`); diff --git a/session.ts b/session.ts index cee91c8..b323144 100644 --- a/session.ts +++ b/session.ts @@ -1,10 +1,10 @@ import { createSocket, RemoteInfo } from "node:dgram"; -import { create_P2pAlive, DevSerial } from "./impl.js"; -import { Commands, CommandsByValue } from "./datatypes.js"; -import { handle_P2PAlive, 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 { Commands, CommandsByValue } from "./datatypes.js"; +import { handle_Drw, handle_P2PAlive, handle_P2PRdy, makePunchPkt, noop, notImpl } from "./handlers.js"; +import { hexdump } from "./hexdump.js"; +import { create_P2pAlive, DevSerial, SendStartVideo, SendVideoResolution, SendWifiDetails } from "./impl.js"; import { opt } from "./options.js"; export type Session = { @@ -17,6 +17,9 @@ export type Session = { connected: boolean; devName: string; timers: ReturnType[]; + curImage: Buffer | null; + rcvSeqId: number; + frame_is_bad: boolean; }; export type PacketHandler = (session: Session, dv: DataView, rinfo: RemoteInfo) => void; @@ -82,7 +85,7 @@ export const makeSession = ( eventEmitter: new EventEmitter(), connected: true, timers: [sessTimer], - devName: dev.serial, + devName: dev.devId, send: (msg: DataView) => { const raw = msg.readU16(); const cmd = CommandsByValue[raw]; @@ -95,6 +98,9 @@ export const makeSession = ( sock.send(new Uint8Array(msg.buffer), SEND_PORT, session.dst_ip); }, dst_ip: ra.address, + curImage: null, + rcvSeqId: 0, + frame_is_bad: false, }; session.eventEmitter.on("disconnect", () => {