fix: prevent agent frame backpressure
This commit is contained in:
+31
-4
@@ -20,6 +20,7 @@ import { initMemory } from './memory/db.js';
|
|||||||
const POD_ROOM = process.env.POD_ROOM ?? 'demo-pod';
|
const POD_ROOM = process.env.POD_ROOM ?? 'demo-pod';
|
||||||
const HERMES_IDENTITY = 'podman-hermes';
|
const HERMES_IDENTITY = 'podman-hermes';
|
||||||
const SAMPLE_INTERVAL_MS = 1000; // ~1 fps to the vision model
|
const SAMPLE_INTERVAL_MS = 1000; // ~1 fps to the vision model
|
||||||
|
const SHUTDOWN_GRACE_MS = 5000;
|
||||||
|
|
||||||
async function agentToken(room: string): Promise<string> {
|
async function agentToken(room: string): Promise<string> {
|
||||||
const at = new AccessToken(env.LIVEKIT_API_KEY, env.LIVEKIT_API_SECRET, {
|
const at = new AccessToken(env.LIVEKIT_API_KEY, env.LIVEKIT_API_SECRET, {
|
||||||
@@ -31,6 +32,20 @@ async function agentToken(room: string): Promise<string> {
|
|||||||
return at.toJwt();
|
return at.toJwt();
|
||||||
}
|
}
|
||||||
|
|
||||||
|
async function withTimeout<T>(promise: Promise<T>, ms: number): Promise<T | null> {
|
||||||
|
let timer: NodeJS.Timeout | undefined;
|
||||||
|
try {
|
||||||
|
return await Promise.race([
|
||||||
|
promise,
|
||||||
|
new Promise<null>((resolve) => {
|
||||||
|
timer = setTimeout(() => resolve(null), ms);
|
||||||
|
}),
|
||||||
|
]);
|
||||||
|
} finally {
|
||||||
|
if (timer) clearTimeout(timer);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
async function main() {
|
async function main() {
|
||||||
// MongoDB is mandatory. Verify the connection before joining the room so bad
|
// MongoDB is mandatory. Verify the connection before joining the room so bad
|
||||||
// creds / unreachable Atlas fail loudly at boot, not silently mid-demo.
|
// creds / unreachable Atlas fail loudly at boot, not silently mid-demo.
|
||||||
@@ -46,6 +61,7 @@ async function main() {
|
|||||||
console.log(`[agent] ${HERMES_IDENTITY} joined room ${POD_ROOM}`);
|
console.log(`[agent] ${HERMES_IDENTITY} joined room ${POD_ROOM}`);
|
||||||
|
|
||||||
const lastSent = new Map<string, number>();
|
const lastSent = new Map<string, number>();
|
||||||
|
const inFlight = new Set<string>();
|
||||||
const activeStreams = new Map<string, ReadableStreamDefaultReader<VideoFrameEvent>>();
|
const activeStreams = new Map<string, ReadableStreamDefaultReader<VideoFrameEvent>>();
|
||||||
|
|
||||||
const streamKey = (
|
const streamKey = (
|
||||||
@@ -69,8 +85,11 @@ async function main() {
|
|||||||
const processFrame = async (engineerId: string, event: VideoFrameEvent) => {
|
const processFrame = async (engineerId: string, event: VideoFrameEvent) => {
|
||||||
const now = Date.now();
|
const now = Date.now();
|
||||||
if (now - (lastSent.get(engineerId) ?? 0) < SAMPLE_INTERVAL_MS) return;
|
if (now - (lastSent.get(engineerId) ?? 0) < SAMPLE_INTERVAL_MS) return;
|
||||||
|
if (inFlight.has(engineerId)) return;
|
||||||
lastSent.set(engineerId, now);
|
lastSent.set(engineerId, now);
|
||||||
|
inFlight.add(engineerId);
|
||||||
|
|
||||||
|
try {
|
||||||
const rgba = event.frame.convert(VideoBufferType.RGBA);
|
const rgba = event.frame.convert(VideoBufferType.RGBA);
|
||||||
const jpeg = await sharp(Buffer.from(rgba.data), {
|
const jpeg = await sharp(Buffer.from(rgba.data), {
|
||||||
raw: { width: rgba.width, height: rgba.height, channels: 4 },
|
raw: { width: rgba.width, height: rgba.height, channels: 4 },
|
||||||
@@ -79,6 +98,9 @@ async function main() {
|
|||||||
.jpeg({ quality: 70 })
|
.jpeg({ quality: 70 })
|
||||||
.toBuffer();
|
.toBuffer();
|
||||||
await podman.onScreenFrame(engineerId, jpeg);
|
await podman.onScreenFrame(engineerId, jpeg);
|
||||||
|
} finally {
|
||||||
|
inFlight.delete(engineerId);
|
||||||
|
}
|
||||||
};
|
};
|
||||||
|
|
||||||
room.on(
|
room.on(
|
||||||
@@ -97,7 +119,7 @@ async function main() {
|
|||||||
while (activeStreams.get(key) === reader) {
|
while (activeStreams.get(key) === reader) {
|
||||||
const { done, value } = await reader.read();
|
const { done, value } = await reader.read();
|
||||||
if (done) break;
|
if (done) break;
|
||||||
await processFrame(id, value).catch((err) =>
|
void processFrame(id, value).catch((err) =>
|
||||||
console.error(`[agent] frame sample failed for ${id}: ${(err as Error).message}`),
|
console.error(`[agent] frame sample failed for ${id}: ${(err as Error).message}`),
|
||||||
);
|
);
|
||||||
}
|
}
|
||||||
@@ -129,10 +151,15 @@ async function main() {
|
|||||||
}
|
}
|
||||||
});
|
});
|
||||||
|
|
||||||
|
let shuttingDown = false;
|
||||||
const shutdown = async () => {
|
const shutdown = async () => {
|
||||||
await Promise.all([...activeStreams.keys()].map(stopStream));
|
if (shuttingDown) return;
|
||||||
await room.disconnect();
|
shuttingDown = true;
|
||||||
await dispose();
|
await withTimeout(
|
||||||
|
Promise.all([...activeStreams.keys()].map(stopStream)).then(() => room.disconnect()),
|
||||||
|
SHUTDOWN_GRACE_MS,
|
||||||
|
);
|
||||||
|
await withTimeout(dispose(), SHUTDOWN_GRACE_MS);
|
||||||
process.exit(0);
|
process.exit(0);
|
||||||
};
|
};
|
||||||
process.on('SIGINT', shutdown);
|
process.on('SIGINT', shutdown);
|
||||||
|
|||||||
Reference in New Issue
Block a user