mirror of
https://github.com/stablyai/orca.git
synced 2026-10-09 08:02:35 +00:00
perf(emulator): stop obsolete video streams after renderer reloads (#22967)
* fix: fence callbacks from retired emulator streams * fix: restore catalog entries required by the current CI baseline * fix(i18n): make AI the recipient of diff notes
This commit is contained in:
@@ -0,0 +1,71 @@
|
||||
import { EventEmitter, once } from 'node:events'
|
||||
import { createServer, type Server } from 'node:http'
|
||||
import type { Socket } from 'node:net'
|
||||
import { afterEach, expect, it, vi } from 'vitest'
|
||||
|
||||
const handlers = new Map<string, (event: unknown, args: unknown) => unknown>()
|
||||
vi.mock('electron', () => ({
|
||||
ipcMain: {
|
||||
handle: (channel: string, handler: (event: unknown, args: unknown) => unknown) => {
|
||||
handlers.set(channel, handler)
|
||||
}
|
||||
},
|
||||
BrowserWindow: { fromWebContents: () => ({}) }
|
||||
}))
|
||||
|
||||
import { registerEmulatorFrameStreamHandlers } from './emulator-frame-stream'
|
||||
|
||||
class Owner extends EventEmitter {
|
||||
isDestroyed = (): boolean => false
|
||||
send = (): void => {
|
||||
this.emit('frame-received')
|
||||
}
|
||||
}
|
||||
|
||||
let server: Server | null = null
|
||||
let owner: Owner | null = null
|
||||
|
||||
afterEach(async () => {
|
||||
owner?.emit('destroyed')
|
||||
owner = null
|
||||
if (server) {
|
||||
server.closeAllConnections()
|
||||
await new Promise<void>((resolve) => server?.close(() => resolve()))
|
||||
server = null
|
||||
}
|
||||
})
|
||||
|
||||
it.each(['did-navigate', 'render-process-gone', 'destroyed'])(
|
||||
'closes the live MJPEG HTTP socket after %s',
|
||||
async (goneEvent) => {
|
||||
registerEmulatorFrameStreamHandlers()
|
||||
const sockets = new Set<Socket>()
|
||||
const httpServer = createServer((_request, response) => {
|
||||
response.writeHead(200, { 'content-type': 'application/octet-stream' })
|
||||
response.write(Buffer.from([0xff, 0xd8, 0x01, 0x02, 0xff, 0xd9]))
|
||||
})
|
||||
server = httpServer
|
||||
httpServer.on('connection', (socket) => {
|
||||
sockets.add(socket)
|
||||
socket.once('close', () => sockets.delete(socket))
|
||||
})
|
||||
await new Promise<void>((resolve) => httpServer.listen(0, '127.0.0.1', resolve))
|
||||
const address = httpServer.address()
|
||||
if (!address || typeof address === 'string') {
|
||||
throw new Error('Expected TCP address')
|
||||
}
|
||||
const sender = new Owner()
|
||||
owner = sender
|
||||
const frameReceived = once(sender, 'frame-received')
|
||||
handlers.get('emulator:frameStreamStart')?.(
|
||||
{ sender },
|
||||
{ streamUrl: `http://127.0.0.1:${address.port}/stream.mjpeg` }
|
||||
)
|
||||
await frameReceived
|
||||
expect(sockets.size).toBe(1)
|
||||
const closed = Promise.all(Array.from(sockets, (socket) => once(socket, 'close')))
|
||||
sender.emit(goneEvent)
|
||||
await closed
|
||||
expect(sockets.size).toBe(0)
|
||||
}
|
||||
)
|
||||
@@ -1,11 +1,11 @@
|
||||
import { BrowserWindow, ipcMain, type WebContents } from 'electron'
|
||||
import { BrowserWindow, ipcMain } from 'electron'
|
||||
import { randomUUID } from 'node:crypto'
|
||||
import { MjpegFrameStream } from '../emulator/mjpeg-frame-stream'
|
||||
import { abortWhenRendererGone } from './renderer-lifetime-abort'
|
||||
|
||||
type FrameStreamSession = {
|
||||
owner: WebContents
|
||||
stream: MjpegFrameStream
|
||||
onOwnerDestroyed: () => void
|
||||
disposeLifetime: () => void
|
||||
}
|
||||
|
||||
const sessions = new Map<string, FrameStreamSession>()
|
||||
@@ -15,11 +15,9 @@ function stopFrameStream(streamId: string): void {
|
||||
if (!session) {
|
||||
return
|
||||
}
|
||||
session.stream.stop()
|
||||
// Why: `.once('destroyed')` self-removes only when that event fires (window
|
||||
// close), so an explicit stop must drop it or each show/hide cycle leaks one.
|
||||
session.owner.removeListener('destroyed', session.onOwnerDestroyed)
|
||||
sessions.delete(streamId)
|
||||
session.disposeLifetime()
|
||||
session.stream.stop()
|
||||
}
|
||||
|
||||
function frameToArrayBuffer(frame: Buffer<ArrayBufferLike>): ArrayBuffer {
|
||||
@@ -45,12 +43,12 @@ export function registerEmulatorFrameStreamHandlers(): void {
|
||||
args.streamUrl,
|
||||
{
|
||||
onError: (message) => {
|
||||
if (!owner.isDestroyed()) {
|
||||
if (sessions.has(streamId) && !owner.isDestroyed()) {
|
||||
owner.send('emulator:frameStreamError', { streamId, message })
|
||||
}
|
||||
},
|
||||
onFrame: (frame) => {
|
||||
if (!owner.isDestroyed()) {
|
||||
if (sessions.has(streamId) && !owner.isDestroyed()) {
|
||||
owner.send('emulator:frameStreamFrame', {
|
||||
streamId,
|
||||
bytes: frameToArrayBuffer(frame)
|
||||
@@ -61,10 +59,22 @@ export function registerEmulatorFrameStreamHandlers(): void {
|
||||
args.streamKey
|
||||
)
|
||||
|
||||
const onOwnerDestroyed = (): void => stopFrameStream(streamId)
|
||||
sessions.set(streamId, { owner, stream, onOwnerDestroyed })
|
||||
owner.once('destroyed', onOwnerDestroyed)
|
||||
stream.start()
|
||||
const lifetime = abortWhenRendererGone(owner)
|
||||
const onRendererGone = (): void => stopFrameStream(streamId)
|
||||
sessions.set(streamId, {
|
||||
stream,
|
||||
disposeLifetime: () => {
|
||||
lifetime.signal.removeEventListener('abort', onRendererGone)
|
||||
lifetime.dispose()
|
||||
}
|
||||
})
|
||||
lifetime.signal.addEventListener('abort', onRendererGone, { once: true })
|
||||
try {
|
||||
stream.start()
|
||||
} catch (error) {
|
||||
stopFrameStream(streamId)
|
||||
throw error
|
||||
}
|
||||
return { streamId }
|
||||
}
|
||||
)
|
||||
|
||||
@@ -0,0 +1,233 @@
|
||||
import { EventEmitter } from 'node:events'
|
||||
import { afterEach, beforeEach, expect, it, vi } from 'vitest'
|
||||
import type { MjpegFrameStreamCallbacks } from '../emulator/mjpeg-frame-stream'
|
||||
import type { ScrcpyVideoSubscriber } from '../emulator/scrcpy-video-registry'
|
||||
|
||||
const mocks = vi.hoisted(() => ({
|
||||
handlers: new Map<string, (event: unknown, args: unknown) => unknown>(),
|
||||
frameCallbacks: new Set<MjpegFrameStreamCallbacks>(),
|
||||
videoSubscribers: new Set<ScrcpyVideoSubscriber>(),
|
||||
frameStarts: vi.fn(),
|
||||
frameStops: vi.fn(),
|
||||
videoStarts: vi.fn(),
|
||||
videoStops: vi.fn()
|
||||
}))
|
||||
|
||||
vi.mock('electron', () => ({
|
||||
ipcMain: {
|
||||
handle: (channel: string, handler: (event: unknown, args: unknown) => unknown) => {
|
||||
mocks.handlers.set(channel, handler)
|
||||
}
|
||||
},
|
||||
BrowserWindow: { fromWebContents: () => ({}) }
|
||||
}))
|
||||
vi.mock('../emulator/mjpeg-frame-stream', () => ({
|
||||
MjpegFrameStream: class {
|
||||
constructor(
|
||||
_url: string,
|
||||
private callbacks: MjpegFrameStreamCallbacks
|
||||
) {}
|
||||
start(): void {
|
||||
mocks.frameStarts()
|
||||
mocks.frameCallbacks.add(this.callbacks)
|
||||
}
|
||||
stop(): void {
|
||||
mocks.frameStops()
|
||||
mocks.frameCallbacks.delete(this.callbacks)
|
||||
}
|
||||
}
|
||||
}))
|
||||
vi.mock('../emulator/scrcpy-video-registry', () => ({
|
||||
scrcpyVideoRegistry: {
|
||||
subscribe: (_deviceId: string, subscriber: ScrcpyVideoSubscriber) => {
|
||||
mocks.videoStarts()
|
||||
mocks.videoSubscribers.add(subscriber)
|
||||
return () => {
|
||||
mocks.videoStops()
|
||||
mocks.videoSubscribers.delete(subscriber)
|
||||
}
|
||||
}
|
||||
}
|
||||
}))
|
||||
vi.mock('../emulator/emulator-probe', () => ({ emulatorProbe: () => {} }))
|
||||
|
||||
import { registerEmulatorFrameStreamHandlers } from './emulator-frame-stream'
|
||||
import { registerEmulatorVideoStreamHandlers } from './emulator-video-stream'
|
||||
|
||||
class Owner extends EventEmitter {
|
||||
send = vi.fn()
|
||||
isDestroyed = (): boolean => false
|
||||
}
|
||||
|
||||
const owners: Owner[] = []
|
||||
const goneEvents = ['did-navigate', 'render-process-gone', 'destroyed'] as const
|
||||
|
||||
function owner(): Owner {
|
||||
const sender = new Owner()
|
||||
owners.push(sender)
|
||||
return sender
|
||||
}
|
||||
|
||||
function start(sender: Owner, kind: 'frame' | 'video'): string {
|
||||
const result = mocks.handlers.get(`emulator:${kind}StreamStart`)?.(
|
||||
{ sender },
|
||||
kind === 'frame'
|
||||
? { streamUrl: 'http://127.0.0.1:0/stream.mjpeg' }
|
||||
: { deviceId: 'emulator-5554' }
|
||||
)
|
||||
if (!result || typeof result !== 'object' || !('streamId' in result)) {
|
||||
throw new Error('Missing stream result')
|
||||
}
|
||||
if (typeof result.streamId !== 'string') {
|
||||
throw new Error('Missing stream id')
|
||||
}
|
||||
return result.streamId
|
||||
}
|
||||
|
||||
function sendFrames(): void {
|
||||
for (const callbacks of mocks.frameCallbacks) {
|
||||
callbacks.onFrame(Buffer.from([0xff, 0xd8, 0xff, 0xd9]))
|
||||
}
|
||||
for (const subscriber of mocks.videoSubscribers) {
|
||||
subscriber({
|
||||
type: 'frame',
|
||||
frame: { config: false, keyFrame: true, pts: '0', bytes: new ArrayBuffer(4) }
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
beforeEach(() => {
|
||||
vi.useFakeTimers()
|
||||
vi.clearAllMocks()
|
||||
registerEmulatorFrameStreamHandlers()
|
||||
registerEmulatorVideoStreamHandlers()
|
||||
})
|
||||
|
||||
afterEach(() => {
|
||||
for (const sender of owners.splice(0)) {
|
||||
sender.emit('destroyed')
|
||||
}
|
||||
vi.clearAllTimers()
|
||||
vi.useRealTimers()
|
||||
mocks.frameCallbacks.clear()
|
||||
mocks.videoSubscribers.clear()
|
||||
})
|
||||
|
||||
it.each(goneEvents)('stops both live streams on %s and leaves no per-frame delivery', (event) => {
|
||||
const sender = owner()
|
||||
start(sender, 'frame')
|
||||
start(sender, 'video')
|
||||
vi.runOnlyPendingTimers()
|
||||
sendFrames()
|
||||
expect(sender.send).toHaveBeenCalledTimes(2)
|
||||
|
||||
sender.emit(event)
|
||||
expect(mocks.frameCallbacks.size).toBe(0)
|
||||
expect(mocks.videoSubscribers.size).toBe(0)
|
||||
expect(mocks.frameStops).toHaveBeenCalledTimes(1)
|
||||
expect(mocks.videoStops).toHaveBeenCalledTimes(1)
|
||||
sendFrames()
|
||||
expect(sender.send).toHaveBeenCalledTimes(2)
|
||||
for (const gone of goneEvents) {
|
||||
expect(sender.listenerCount(gone)).toBe(0)
|
||||
}
|
||||
})
|
||||
|
||||
it.each(goneEvents)('does not start deferred video work after %s', (event) => {
|
||||
const sender = owner()
|
||||
start(sender, 'video')
|
||||
sender.emit(event)
|
||||
vi.runOnlyPendingTimers()
|
||||
expect(mocks.videoStarts).not.toHaveBeenCalled()
|
||||
expect(mocks.videoSubscribers.size).toBe(0)
|
||||
})
|
||||
|
||||
it('keeps only the current document streams after repeated reloads', () => {
|
||||
const sender = owner()
|
||||
for (let cycle = 0; cycle < 15; cycle++) {
|
||||
start(sender, 'frame')
|
||||
start(sender, 'video')
|
||||
vi.runOnlyPendingTimers()
|
||||
sender.emit('did-navigate')
|
||||
}
|
||||
start(sender, 'frame')
|
||||
start(sender, 'video')
|
||||
vi.runOnlyPendingTimers()
|
||||
expect(mocks.frameCallbacks.size).toBe(1)
|
||||
expect(mocks.videoSubscribers.size).toBe(1)
|
||||
sendFrames()
|
||||
expect(sender.send).toHaveBeenCalledTimes(2)
|
||||
})
|
||||
|
||||
it('keeps streams through same-document or prevented navigation and releases on explicit stop', () => {
|
||||
const sender = owner()
|
||||
const frameId = start(sender, 'frame')
|
||||
const videoId = start(sender, 'video')
|
||||
vi.runOnlyPendingTimers()
|
||||
sender.emit('did-start-navigation')
|
||||
sender.emit('did-navigate-in-page')
|
||||
sendFrames()
|
||||
expect(sender.send).toHaveBeenCalledTimes(2)
|
||||
mocks.handlers.get('emulator:frameStreamStop')?.({ sender }, { streamId: frameId })
|
||||
mocks.handlers.get('emulator:videoStreamStop')?.({ sender }, { streamId: videoId })
|
||||
expect(mocks.frameCallbacks.size).toBe(0)
|
||||
expect(mocks.videoSubscribers.size).toBe(0)
|
||||
for (const event of goneEvents) {
|
||||
expect(sender.listenerCount(event)).toBe(0)
|
||||
}
|
||||
})
|
||||
|
||||
it('releases a failed frame-stream start without retaining renderer listeners', () => {
|
||||
const sender = owner()
|
||||
mocks.frameStarts.mockImplementationOnce(() => {
|
||||
throw new Error('Stream unavailable')
|
||||
})
|
||||
expect(() => start(sender, 'frame')).toThrow('Stream unavailable')
|
||||
expect(mocks.frameStops).toHaveBeenCalledTimes(1)
|
||||
for (const event of goneEvents) {
|
||||
expect(sender.listenerCount(event)).toBe(0)
|
||||
}
|
||||
})
|
||||
|
||||
it('keeps another renderer streams live when the first document reloads', () => {
|
||||
const first = owner()
|
||||
const second = owner()
|
||||
for (const sender of [first, second]) {
|
||||
start(sender, 'frame')
|
||||
start(sender, 'video')
|
||||
}
|
||||
vi.runOnlyPendingTimers()
|
||||
first.emit('did-navigate')
|
||||
sendFrames()
|
||||
expect(first.send).not.toHaveBeenCalled()
|
||||
expect(second.send).toHaveBeenCalledTimes(2)
|
||||
expect(mocks.frameCallbacks.size).toBe(1)
|
||||
expect(mocks.videoSubscribers.size).toBe(1)
|
||||
})
|
||||
|
||||
it.each([...goneEvents, 'explicit stop'] as const)(
|
||||
'suppresses retired frame callbacks after %s while keeping current stream errors',
|
||||
(event) => {
|
||||
const sender = owner()
|
||||
const streamId = start(sender, 'frame')
|
||||
const callbacks = [...mocks.frameCallbacks][0]
|
||||
callbacks.onError('active failure')
|
||||
expect(sender.send).toHaveBeenCalledWith('emulator:frameStreamError', {
|
||||
streamId,
|
||||
message: 'active failure'
|
||||
})
|
||||
sender.send.mockClear()
|
||||
mocks.frameStops.mockImplementationOnce(() => callbacks.onError('stop failure'))
|
||||
if (event === 'explicit stop') {
|
||||
mocks.handlers.get('emulator:frameStreamStop')?.({ sender }, { streamId })
|
||||
} else {
|
||||
sender.emit(event)
|
||||
}
|
||||
start(sender, 'frame')
|
||||
callbacks.onError('late failure')
|
||||
callbacks.onFrame(Buffer.from([0xff, 0xd8, 0xff, 0xd9]))
|
||||
expect(sender.send).not.toHaveBeenCalled()
|
||||
sendFrames()
|
||||
expect(sender.send).toHaveBeenCalledOnce()
|
||||
}
|
||||
)
|
||||
@@ -2,6 +2,7 @@ import { BrowserWindow, ipcMain, type WebContents } from 'electron'
|
||||
import { randomUUID } from 'node:crypto'
|
||||
import { scrcpyVideoRegistry } from '../emulator/scrcpy-video-registry'
|
||||
import { emulatorProbe } from '../emulator/emulator-probe'
|
||||
import { abortWhenRendererGone } from './renderer-lifetime-abort'
|
||||
|
||||
// Bridges the main-process scrcpy video registry to renderer subscribers. The
|
||||
// renderer calls emulator:videoStreamStart with a deviceId; meta + H.264 access
|
||||
@@ -11,7 +12,8 @@ export function registerEmulatorVideoStreamHandlers(): void {
|
||||
type Subscription = {
|
||||
owner: WebContents
|
||||
unsubscribe: () => void
|
||||
onOwnerDestroyed: () => void
|
||||
disposeLifetime: () => void
|
||||
startTimer: ReturnType<typeof setTimeout> | null
|
||||
}
|
||||
const subscriptions = new Map<string, Subscription>()
|
||||
|
||||
@@ -20,11 +22,12 @@ export function registerEmulatorVideoStreamHandlers(): void {
|
||||
if (!subscription || (owner && subscription.owner !== owner)) {
|
||||
return
|
||||
}
|
||||
subscription.unsubscribe()
|
||||
// Why: `.once('destroyed')` self-removes only when that event fires (window
|
||||
// close), so an explicit stop must drop it or each show/hide cycle leaks one.
|
||||
subscription.owner.removeListener('destroyed', subscription.onOwnerDestroyed)
|
||||
subscriptions.delete(streamId)
|
||||
if (subscription.startTimer !== null) {
|
||||
clearTimeout(subscription.startTimer)
|
||||
}
|
||||
subscription.disposeLifetime()
|
||||
subscription.unsubscribe()
|
||||
}
|
||||
|
||||
ipcMain.handle(
|
||||
@@ -44,14 +47,21 @@ export function registerEmulatorVideoStreamHandlers(): void {
|
||||
throw new Error('Video stream id is already in use by another renderer')
|
||||
}
|
||||
stopSubscription(streamId, owner)
|
||||
const onOwnerDestroyed = (): void => stopSubscription(streamId, owner)
|
||||
const lifetime = abortWhenRendererGone(owner)
|
||||
const onRendererGone = (): void => stopSubscription(streamId, owner)
|
||||
const pendingSubscription: Subscription = {
|
||||
owner,
|
||||
unsubscribe: () => {},
|
||||
onOwnerDestroyed
|
||||
disposeLifetime: () => {
|
||||
lifetime.signal.removeEventListener('abort', onRendererGone)
|
||||
lifetime.dispose()
|
||||
},
|
||||
startTimer: null
|
||||
}
|
||||
subscriptions.set(streamId, pendingSubscription)
|
||||
setTimeout(() => {
|
||||
lifetime.signal.addEventListener('abort', onRendererGone, { once: true })
|
||||
pendingSubscription.startTimer = setTimeout(() => {
|
||||
pendingSubscription.startTimer = null
|
||||
if (owner.isDestroyed() || subscriptions.get(streamId) !== pendingSubscription) {
|
||||
return
|
||||
}
|
||||
@@ -75,7 +85,6 @@ export function registerEmulatorVideoStreamHandlers(): void {
|
||||
})
|
||||
pendingSubscription.unsubscribe = unsubscribe
|
||||
}, 0)
|
||||
owner.once('destroyed', onOwnerDestroyed)
|
||||
return { streamId }
|
||||
}
|
||||
)
|
||||
|
||||
Reference in New Issue
Block a user