support multiple streams concurrently
This commit is contained in:
+13
-16
@@ -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)]);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
+22
-13
@@ -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<string, ServerResponse[]> = {};
|
||||
let sessions: Record<string, Session> = {};
|
||||
|
||||
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(`<html>`);
|
||||
Object.keys(sessions).forEach((id) => res.write(`<a href="/camera/${id}">${id}</a>`));
|
||||
Object.keys(sessions).forEach((id) => res.write(`<a href="/camera/${id}"><img src="/camera/${id}"/></a><hr/>`));
|
||||
res.write(`</html>`);
|
||||
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`);
|
||||
|
||||
+12
-6
@@ -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<typeof setInterval>[];
|
||||
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", () => {
|
||||
|
||||
Reference in New Issue
Block a user