From 1981b7d6a9417a96e3099e99a30933b7360701c7 Mon Sep 17 00:00:00 2001 From: harveysang Date: Thu, 17 Sep 2026 23:53:24 +0800 Subject: [PATCH] Share scoped session cleanup and video transport across host instances --- .../device-runtime/src/forward-registry.ts | 280 +++++++++ packages/device-runtime/src/scrcpy-stream.ts | 546 +++++++++++++++++ .../device-runtime/src/session-runtime.ts | 80 +++ .../device-runtime/tests/session-contract.ts | 64 ++ plugins/opengui/src/codex/service.ts | 55 +- plugins/opengui/src/daemon.ts | 22 +- plugins/opengui/src/forward-registry.ts | 206 +------ plugins/opengui/src/scrcpy-stream.ts | 554 +----------------- plugins/opengui/tests/daemon.spec.ts | 12 + .../opengui/tests/runtime-contract.spec.ts | 2 + workbuddy-plugin/src/forward-registry.ts | 289 +-------- workbuddy-plugin/src/scrcpy-stream.ts | 554 +----------------- workbuddy-plugin/src/service.ts | 56 +- .../tests/runtime-contract.spec.ts | 2 + 14 files changed, 1069 insertions(+), 1653 deletions(-) create mode 100644 packages/device-runtime/src/forward-registry.ts create mode 100644 packages/device-runtime/src/scrcpy-stream.ts create mode 100644 packages/device-runtime/src/session-runtime.ts create mode 100644 packages/device-runtime/tests/session-contract.ts diff --git a/packages/device-runtime/src/forward-registry.ts b/packages/device-runtime/src/forward-registry.ts new file mode 100644 index 0000000..0e1494c --- /dev/null +++ b/packages/device-runtime/src/forward-registry.ts @@ -0,0 +1,280 @@ +import { randomUUID } from 'node:crypto' +import { mkdir, open, readFile, rename, rm, stat, writeFile } from 'node:fs/promises' +import { closeSync, openSync, readFileSync, renameSync, rmSync, writeFileSync } from 'node:fs' +import { spawnSync } from 'node:child_process' +import { dirname } from 'node:path' + +export type OwnedForwardKind = 'text-input' | 'video-stream' + +export interface OwnedForward { + readonly serial: string + readonly port: number + readonly scid: string + readonly kind: OwnedForwardKind +} + +interface StoredForward extends OwnedForward { + readonly ownerId?: string + readonly ownerPid?: number +} + +export type ForwardAdbRunner = (args: readonly string[], signal: AbortSignal) => Promise + +interface ListedForward { + readonly serial: string + readonly local: string + readonly remote: string +} + +export function parseAdbForwardList(output: string): ListedForward[] { + const forwards: ListedForward[] = [] + for (const line of output.split(/\r?\n/u)) { + const [serial, local, remote, ...extra] = line.trim().split(/\s+/u) + if (!serial || !local || !remote || extra.length > 0) continue + forwards.push({ serial, local, remote }) + } + return forwards +} + +/** Durable inventory of ADB forwards created by this plugin, never inferred from global ADB state. */ +export class OwnedForwardRegistry { + private mutation = Promise.resolve() + private readonly ownerId = randomUUID() + + constructor(private readonly path: string) {} + + async track(record: OwnedForward): Promise { + await this.mutate(records => { + records.set(this.key(record), { ...record, ownerId: this.ownerId, ownerPid: process.pid }) + }) + } + + async release(record: OwnedForward, runAdb: ForwardAdbRunner, timeoutMs = 5_000): Promise { + let forwards: ListedForward[] + try { + const listed = await runAdb( + ['-s', record.serial, 'forward', '--list'], + AbortSignal.timeout(timeoutMs), + ) + forwards = parseAdbForwardList(String(listed ?? '')) + } catch { + return false + } + const local = `tcp:${record.port}` + const owned = forwards.find(candidate => candidate.serial === record.serial && candidate.local === local) + if (owned === undefined || owned.remote !== this.remote(record)) { + try { + await this.deleteMatching(record) + return true + } catch { + return false + } + } + try { + await runAdb( + ['-s', record.serial, 'forward', '--remove', local], + AbortSignal.timeout(timeoutMs), + ) + } catch { + return false + } + try { + await this.deleteMatching(record) + return true + } catch { + return false + } + } + + async recover(runAdb: ForwardAdbRunner): Promise<{ removed: number; retained: number }> { + const records = await this.read() + let removed = 0 + let retained = 0 + for (const record of records.values()) { + if (record.ownerId !== undefined && record.ownerId !== this.ownerId + && record.ownerPid !== undefined && this.processAlive(record.ownerPid)) { + retained += 1 + continue + } + if (await this.release(record, runAdb)) removed += 1 + else retained += 1 + } + return { removed, retained } + } + + async list(): Promise { + return [...(await this.read()).values()].map(record => this.publicRecord(record)) + } + + /** Best-effort synchronous signal-path cleanup when the Host cannot await plugin disposal. */ + releaseAllSync(adbPath: string): { removed: number; retained: number } { + const lock = `${this.path}.lock` + let lockFd: number + try { + lockFd = openSync(lock, 'wx', 0o600) + } catch { + return { removed: 0, retained: this.readSync().filter(record => record.ownerId === this.ownerId).length } + } + let records: StoredForward[] + try { + const value = JSON.parse(readFileSync(this.path, 'utf8')) as unknown + records = Array.isArray(value) ? value.filter(candidate => this.valid(candidate)) : [] + } catch { + closeSync(lockFd) + rmSync(lock, { force: true }) + return { removed: 0, retained: 0 } + } + const ownedRecords = records.filter(record => record.ownerId === this.ownerId) + const retainedOwned: StoredForward[] = [] + for (const record of ownedRecords) { + const listed = spawnSync(adbPath, ['-s', record.serial, 'forward', '--list'], { + timeout: 1_500, + encoding: 'utf8', + stdio: ['ignore', 'pipe', 'ignore'], + }) + if (listed.status !== 0 || listed.error !== undefined) { + retainedOwned.push(record) + continue + } + const local = `tcp:${record.port}` + const owned = parseAdbForwardList(listed.stdout ?? '') + .find(candidate => candidate.serial === record.serial && candidate.local === local) + if (owned === undefined || owned.remote !== this.remote(record)) continue + const removed = spawnSync(adbPath, ['-s', record.serial, 'forward', '--remove', local], { + timeout: 1_500, + stdio: 'ignore', + }) + if (removed.status !== 0 || removed.error !== undefined) retainedOwned.push(record) + } + const retainedKeys = new Set(retainedOwned.map(record => `${this.key(record)}\u0000${record.scid}`)) + const retained = records.filter(record => record.ownerId !== this.ownerId + || retainedKeys.has(`${this.key(record)}\u0000${record.scid}`)) + try { + if (retained.length === 0) rmSync(this.path, { force: true }) + else { + const temporary = `${this.path}.${process.pid}.${this.ownerId}.signal.tmp` + writeFileSync(temporary, JSON.stringify(retained), { encoding: 'utf8', mode: 0o600 }) + renameSync(temporary, this.path) + } + } catch { /* startup recovery retains the durable fallback */ } + finally { + closeSync(lockFd) + rmSync(lock, { force: true }) + } + return { removed: ownedRecords.length - retainedOwned.length, retained: retainedOwned.length } + } + + private key(record: OwnedForward): string { + return `${record.serial}\u0000${record.port}` + } + + private remote(record: OwnedForward): string { + return `localabstract:scrcpy_${record.scid}` + } + + private valid(value: unknown): value is StoredForward { + if (typeof value !== 'object' || value === null) return false + const candidate = value as Partial + return typeof candidate.serial === 'string' && Number.isSafeInteger(candidate.port) + && typeof candidate.scid === 'string' + && (candidate.kind === 'text-input' || candidate.kind === 'video-stream') + } + + private async mutate(change: (records: Map) => void | boolean): Promise { + const operation = this.mutation.then(async () => { + await mkdir(dirname(this.path), { recursive: true }) + const releaseLock = await this.acquireLock() + try { + const records = await this.read() + change(records) + if (records.size === 0) { + await rm(this.path, { force: true }) + return + } + const temporary = `${this.path}.${process.pid}.${this.ownerId}.${randomUUID()}.tmp` + await writeFile(temporary, JSON.stringify([...records.values()]), { encoding: 'utf8', mode: 0o600 }) + await rename(temporary, this.path) + } finally { + await releaseLock() + } + }) + this.mutation = operation.catch(() => undefined) + await operation + } + + private async read(): Promise> { + try { + const parsed = JSON.parse(await readFile(this.path, 'utf8')) as unknown + if (!Array.isArray(parsed)) return new Map() + const records = new Map() + for (const value of parsed) { + if (!this.valid(value)) continue + const record = value + records.set(this.key(record), record) + } + return records + } catch { + return new Map() + } + } + + private async deleteMatching(record: OwnedForward): Promise { + await this.mutate(records => { + const stored = records.get(this.key(record)) + if (stored?.scid === record.scid) records.delete(this.key(record)) + }) + } + + private publicRecord(record: StoredForward): OwnedForward { + return { serial: record.serial, port: record.port, scid: record.scid, kind: record.kind } + } + + private processAlive(pid: number): boolean { + try { + process.kill(pid, 0) + return true + } catch { + return false + } + } + + private readSync(): StoredForward[] { + try { + const parsed = JSON.parse(readFileSync(this.path, 'utf8')) as unknown + return Array.isArray(parsed) ? parsed.filter(candidate => this.valid(candidate)) : [] + } catch { + return [] + } + } + + private async acquireLock(): Promise<() => Promise> { + const lock = `${this.path}.lock` + const deadline = Date.now() + 5_000 + while (true) { + try { + const handle = await open(lock, 'wx', 0o600) + return async () => { + await handle.close() + await rm(lock, { force: true }) + } + } catch (error) { + // Windows can report a sharing violation while another writer deletes its lock. + // Retry acquisition only; never infer that EPERM grants permission to remove it. + if ((error as NodeJS.ErrnoException).code === 'EPERM') { + if (Date.now() >= deadline) throw error + await new Promise(resolve => setTimeout(resolve, 10)) + continue + } + if ((error as NodeJS.ErrnoException).code !== 'EEXIST') throw error + try { + if (Date.now() - (await stat(lock)).mtimeMs > 10_000) { + await rm(lock, { force: true }) + continue + } + } catch { /* another writer released the lock */ } + if (Date.now() >= deadline) throw new Error('opengui: timed out waiting for forward registry lock') + await new Promise(resolve => setTimeout(resolve, 10)) + } + } + } +} diff --git a/packages/device-runtime/src/scrcpy-stream.ts b/packages/device-runtime/src/scrcpy-stream.ts new file mode 100644 index 0000000..ebc2529 --- /dev/null +++ b/packages/device-runtime/src/scrcpy-stream.ts @@ -0,0 +1,546 @@ +// Ported from the repository video transport; see VIDEO-NOTICE.md. +import { randomBytes } from 'node:crypto' +import { spawn } from 'node:child_process' +import type { ChildProcess } from 'node:child_process' +import { connect, createServer } from 'node:net' +import type { Socket } from 'node:net' +import type { VideoDevice, ScrcpyStreamSink } from './contracts.ts' +export type { VideoDevice, ScrcpyStreamSink } from './contracts.ts' +import { OwnedForwardRegistry, type OwnedForward } from './forward-registry.ts' + +export interface ScrcpyAsset { + key: string + archive: 'tar.gz' | 'zip' + archiveRoot: string + executable: string + url: string + bytes: number + sha256: string +} +export interface VideoInstaller { + isInstalled(asset: ScrcpyAsset): Promise + ensure(asset: ScrcpyAsset, signal: AbortSignal, progress: (value: { + phase: 'downloading' | 'extracting'; downloadedBytes?: number; totalBytes?: number + }) => void): Promise<{ server: string }> +} + +const SESSION_PACKET_FLAG = 0x8000000000000000n +const CONFIG_PACKET_FLAG = 0x4000000000000000n +const KEY_PACKET_FLAG = 0x2000000000000000n +const PTS_MASK = 0x1fffffffffffffffn + +export type ScrcpyVideoEvent = { + readonly type: 'codec' + readonly codec: 'h264' +} | { + readonly type: 'session' + readonly width: number + readonly height: number + readonly clientResized: boolean +} | { + readonly type: 'packet' + readonly config: boolean + readonly key: boolean + readonly pts: bigint + readonly data: Buffer +} + +/** Incremental parser for scrcpy 4.1 stream metadata and H.264 media packets. */ +export class ScrcpyVideoPacketParser { + private buffer = Buffer.alloc(0) + private codecRead = false + + push(chunk: Buffer): ScrcpyVideoEvent[] { + if (chunk.length > 0) this.buffer = Buffer.concat([this.buffer, chunk]) + const events: ScrcpyVideoEvent[] = [] + if (!this.codecRead) { + if (this.buffer.length < 4) return events + const codec = this.buffer.subarray(0, 4).toString('ascii') + if (codec !== 'h264') throw new Error(`opengui-workbuddy: unsupported scrcpy video codec ${codec}`) + this.codecRead = true + this.buffer = this.buffer.subarray(4) + events.push({ type: 'codec', codec: 'h264' }) + } + while (this.buffer.length >= 12) { + const flagsAndPts = this.buffer.readBigUInt64BE(0) + if ((flagsAndPts & SESSION_PACKET_FLAG) !== 0n) { + const flags = this.buffer.readUInt32BE(0) + const width = this.buffer.readUInt32BE(4) + const height = this.buffer.readUInt32BE(8) + if (width < 1 || height < 1 || width > 16_384 || height > 16_384) { + throw new Error(`opengui-workbuddy: invalid scrcpy video size ${width}x${height}`) + } + this.buffer = this.buffer.subarray(12) + events.push({ type: 'session', width, height, clientResized: (flags & 1) === 1 }) + continue + } + const size = this.buffer.readUInt32BE(8) + if (size > 16 * 1024 * 1024) throw new Error('opengui-workbuddy: scrcpy video packet exceeds 16 MiB') + if (this.buffer.length < 12 + size) break + const data = Buffer.from(this.buffer.subarray(12, 12 + size)) + this.buffer = this.buffer.subarray(12 + size) + events.push({ + type: 'packet', + config: (flagsAndPts & CONFIG_PACKET_FLAG) !== 0n, + key: (flagsAndPts & KEY_PACKET_FLAG) !== 0n, + pts: flagsAndPts & PTS_MASK, + data, + }) + } + return events + } +} + +/** Fixed read-only server options for the embedded low-latency stream. */ +export function buildScrcpyVideoServerArgs(scid: string, serverPath: string, version: string): string[] { + return [ + `CLASSPATH=${serverPath}`, + 'app_process', '/', 'com.genymobile.scrcpy.Server', version, + `scid=${scid}`, + 'tunnel_forward=true', + 'video=true', + 'audio=false', + 'control=false', + 'cleanup=false', + 'video_codec=h264', + 'max_size=960', + 'max_fps=30', + 'video_bit_rate=2000000', + 'video_codec_options=i-frame-interval=1', + 'send_dummy_byte=false', + 'send_device_meta=false', + 'send_stream_meta=true', + 'send_frame_meta=true', + ] +} + +export interface ScrcpyStreamStatus { + supported: boolean + cached: boolean + /** @deprecated Kept true on supported Hosts for one client-compatibility release. */ + approved: boolean + phase: 'idle' | 'downloading' | 'extracting' | 'ready' | 'error' + version: string + totalBytes?: number + downloadedBytes?: number + activeSources: number + maxSources: number + message?: string +} + +type AdbRunner = (args: readonly string[], signal: AbortSignal) => Promise + +interface StreamEntry { + readonly device: VideoDevice + readonly subscribers: Set + readonly waiting: Set + readonly controller: AbortController + operation: Promise + socket?: Socket + process?: ChildProcess + port?: number + forward?: OwnedForward + idleTimer: ReturnType | undefined + lastCodec?: string + lastSession?: string + replay: Buffer[] + replayBytes: number + closeCode: number + closeReason: string + closed: boolean + closing?: Promise +} + +export interface ScrcpyVideoStreamsOptions { + version: string + remoteServer: string + adbPath: () => string + runAdb: AdbRunner + installer: VideoInstaller + asset?: ScrcpyAsset | undefined + spawn?: typeof spawn + connect?: typeof connect + freePort?: () => Promise + idleGraceMs?: number + maxSources?: number + onError?: (error: unknown) => void + forwardRegistry: OwnedForwardRegistry +} + +/** Shares one scrcpy encoder per device across same-origin browser subscribers. */ +export class ScrcpyVideoStreams { + private readonly installer: VideoInstaller + private readonly asset: ScrcpyAsset | undefined + private readonly spawnImpl: typeof spawn + private readonly connectImpl: typeof connect + private readonly freePort: () => Promise + private readonly idleGraceMs: number + private readonly maxSources: number + private readonly onError: (error: unknown) => void + private readonly forwardRegistry: OwnedForwardRegistry + private readonly entries = new Map() + private readonly lifetime = new AbortController() + private phase: ScrcpyStreamStatus['phase'] = 'idle' + private downloadedBytes: number | undefined + private message: string | undefined + + constructor(private readonly options: ScrcpyVideoStreamsOptions) { + this.installer = options.installer + this.asset = options.asset + this.spawnImpl = options.spawn ?? spawn + this.connectImpl = options.connect ?? connect + this.freePort = options.freePort ?? availableTcpPort + this.idleGraceMs = options.idleGraceMs ?? 2_000 + this.maxSources = options.maxSources ?? 4 + this.onError = options.onError ?? (() => {}) + this.forwardRegistry = options.forwardRegistry + } + + async prepare(signal: AbortSignal): Promise { + if (!this.asset) throw new Error('stream_unsupported') + await this.installer.ensure(this.asset, signal, () => {}) + } + + approve(): boolean { + // Compatibility endpoint: first-use preparation is automatic now. + return this.asset !== undefined + } + + async status(): Promise { + const cached = this.asset !== undefined && await this.installer.isInstalled(this.asset) + if (cached && this.phase === 'idle') this.phase = 'ready' + return { + supported: this.asset !== undefined, + cached, + approved: this.asset !== undefined, + phase: cached && this.phase === 'idle' ? 'ready' : this.phase, + version: this.options.version, + ...(this.asset === undefined ? {} : { totalBytes: this.asset.bytes }), + ...(this.downloadedBytes === undefined ? {} : { downloadedBytes: this.downloadedBytes }), + activeSources: this.entries.size, + maxSources: this.maxSources, + ...(this.message === undefined ? {} : { message: this.message }), + } + } + + async subscribe(device: VideoDevice, sink: ScrcpyStreamSink): Promise<() => void> { + if (this.lifetime.signal.aborted) throw new Error('stream_disposed') + const asset = this.asset + if (asset === undefined) throw new Error('stream_unsupported') + let entry = this.entries.get(device.id) + if (entry?.closed === true) { + // A closing encoder still consumes a source slot until its owned resources drain. + await entry.closing + return this.subscribe(device, sink) + } + if (entry === undefined) { + if (this.entries.size >= this.maxSources) throw new Error('stream_capacity_wait') + const controller = new AbortController() + entry = { + device, + subscribers: new Set(), + waiting: new Set(), + controller, + operation: Promise.resolve(), + idleTimer: undefined, + replay: [], + replayBytes: 0, + closeCode: 1000, + closeReason: 'stream stopped', + closed: false, + } + this.entries.set(device.id, entry) + entry.operation = this.start(entry, asset).catch(error => { + if (!entry!.controller.signal.aborted) { + this.phase = 'error' + this.message = error instanceof Error ? error.message : String(error) + this.onError(error) + this.broadcastText(entry!, { type: 'error', message: this.publicError(error) }) + entry!.closeCode = 1011 + entry!.closeReason = 'stream failed' + } + }).finally(() => { + void this.closeEntry(entry!) + }) + } + if (entry.idleTimer !== undefined) { + clearTimeout(entry.idleTimer) + entry.idleTimer = undefined + } + entry.subscribers.add(sink) + if (entry.lastCodec !== undefined) sink.sendText(entry.lastCodec) + if (entry.lastSession !== undefined) sink.sendText(entry.lastSession) + if (entry.replayBytes <= 1_000_000) { + for (const frame of entry.replay) sink.sendBinary(frame) + } else entry.waiting.add(sink) + return () => this.unsubscribe(entry!, sink) + } + + async dispose(): Promise { + if (!this.lifetime.signal.aborted) this.lifetime.abort(new Error('opengui: stream manager disposed')) + await Promise.allSettled([...this.entries.values()].map(entry => this.closeEntry(entry))) + } + + private unsubscribe(entry: StreamEntry, sink: ScrcpyStreamSink): void { + entry.subscribers.delete(sink) + entry.waiting.delete(sink) + if (entry.subscribers.size > 0 || entry.closed || entry.idleTimer !== undefined) return + entry.idleTimer = setTimeout(() => { + entry.idleTimer = undefined + if (entry.subscribers.size === 0) void this.closeEntry(entry) + }, this.idleGraceMs) + } + + private async start(entry: StreamEntry, asset: ScrcpyAsset): Promise { + const signal = AbortSignal.any([entry.controller.signal, this.lifetime.signal]) + this.phase = 'downloading' + this.message = undefined + const installed = await this.installer.ensure(asset, signal, progress => { + this.phase = progress.phase + this.downloadedBytes = progress.downloadedBytes + this.broadcastText(entry, { type: 'install', phase: progress.phase, downloadedBytes: progress.downloadedBytes, totalBytes: progress.totalBytes }) + }) + this.phase = 'ready' + this.downloadedBytes = undefined + signal.throwIfAborted() + + const port = await this.freePort() + entry.port = port + const scid = (randomBytes(4).readUInt32BE(0) & 0x7fffffff).toString(16).padStart(8, '0') + const forward: OwnedForward = { serial: entry.device.serial, port, scid, kind: 'video-stream' } + await this.options.runAdb(['-s', entry.device.serial, 'push', installed.server, this.options.remoteServer], signal) + try { + await this.forwardRegistry.track(forward) + entry.forward = forward + await this.options.runAdb([ + '-s', entry.device.serial, 'forward', '--no-rebind', `tcp:${port}`, `localabstract:scrcpy_${scid}`, + ], signal) + } catch (error) { + await this.forwardRegistry.release(forward, this.options.runAdb).catch(() => false) + throw error + } + + const child = this.spawnImpl(this.options.adbPath(), [ + '-s', entry.device.serial, 'shell', ...buildScrcpyVideoServerArgs(scid, this.options.remoteServer, this.options.version), + ], { shell: false, windowsHide: true, stdio: ['ignore', 'ignore', 'pipe'] }) + entry.process = child + let stderr = '' + child.stderr?.on('data', (chunk: Buffer | string) => { stderr = `${stderr}${String(chunk)}`.slice(-2_000) }) + await waitForSpawn(child, signal) + const socket = await connectVideo(this.connectImpl, port, signal) + entry.socket = socket + const parser = new ScrcpyVideoPacketParser() + socket.on('data', chunk => { + try { + for (const event of parser.push(Buffer.from(chunk))) this.broadcastEvent(entry, event) + } catch (error) { + this.broadcastText(entry, { type: 'error', message: this.publicError(error) }) + void this.closeEntry(entry) + } + }) + const settled = new Promise((resolve, reject) => { + let socketCloseTimer: ReturnType | undefined + const rejectSocketClose = (): void => { + socketCloseTimer = setTimeout(() => { + reject(new Error(`scrcpy video socket closed unexpectedly${stderr.trim() ? `: ${stderr.trim()}` : ''}`)) + }, 150) + } + socket.once('error', reject) + socket.once('close', () => signal.aborted ? resolve() : rejectSocketClose()) + child.once('exit', (code, exitSignal) => { + if (socketCloseTimer !== undefined) clearTimeout(socketCloseTimer) + if (entry.controller.signal.aborted) resolve() + else reject(new Error(`scrcpy video server exited (code=${String(code)}, signal=${String(exitSignal)})${stderr.trim() ? `: ${stderr.trim()}` : ''}`)) + }) + signal.addEventListener('abort', () => { + if (socketCloseTimer !== undefined) clearTimeout(socketCloseTimer) + resolve() + }, { once: true }) + }) + await settled + } + + private broadcastEvent(entry: StreamEntry, event: ScrcpyVideoEvent): void { + if (event.type === 'packet') { + const frame = this.packetFrame(event) + if (event.config) { + entry.replay = [frame] + entry.replayBytes = frame.byteLength + } else if (event.key) { + entry.replay = [...entry.replay.filter(packet => (packet[0]! & 1) !== 0), frame] + entry.replayBytes = entry.replay.reduce((total, packet) => total + packet.byteLength, 0) + } else if (entry.replay.some(packet => (packet[0]! & 2) !== 0)) { + if (entry.replayBytes + frame.byteLength <= 8 * 1024 * 1024) { + entry.replay.push(frame) + entry.replayBytes += frame.byteLength + } else { + entry.replay = entry.replay.filter(packet => (packet[0]! & 1) !== 0) + entry.replayBytes = entry.replay.reduce((total, packet) => total + packet.byteLength, 0) + } + } + for (const sink of entry.subscribers) { + if (sink.bufferedBytes() > 1_000_000) { + entry.waiting.add(sink) + if (sink.bufferedBytes() > 2_000_000) sink.close(1013, 'slow_client') + continue + } + if (entry.waiting.has(sink)) { + if (!event.key) continue + sink.sendText(JSON.stringify({ type: 'reset' })) + for (const config of entry.replay.filter(packet => (packet[0]! & 1) !== 0)) sink.sendBinary(config) + entry.waiting.delete(sink) + } + sink.sendBinary(frame) + } + return + } + const text = JSON.stringify(event) + if (event.type === 'codec') entry.lastCodec = text + else { + entry.lastSession = text + entry.replay = [] + entry.replayBytes = 0 + } + for (const sink of entry.subscribers) sink.sendText(text) + } + + private packetFrame(event: Extract): Buffer { + const frame = Buffer.allocUnsafe(9 + event.data.length) + frame[0] = (event.config ? 1 : 0) | (event.key ? 2 : 0) + frame.writeBigUInt64BE(event.pts, 1) + event.data.copy(frame, 9) + return frame + } + + private broadcastText(entry: StreamEntry, value: unknown): void { + const text = JSON.stringify(value) + for (const sink of entry.subscribers) sink.sendText(text) + } + + private closeEntry(entry: StreamEntry): Promise { + if (entry.closing !== undefined) return entry.closing + if (entry.closed) return Promise.resolve() + entry.closed = true + const closing = (async (): Promise => { + if (entry.idleTimer !== undefined) clearTimeout(entry.idleTimer) + entry.controller.abort(new Error('opengui-workbuddy: embedded stream stopped')) + entry.socket?.destroy() + const child = entry.process + if (child !== undefined && child.exitCode === null) await terminateChild(child) + if (entry.forward !== undefined) await this.forwardRegistry.release(entry.forward, this.options.runAdb, 1_000) + if (this.entries.get(entry.device.id) === entry) this.entries.delete(entry.device.id) + for (const sink of entry.subscribers) sink.close(entry.closeCode, entry.closeReason) + entry.subscribers.clear() + })() + entry.closing = closing + return closing + } + + private publicError(_error: unknown): string { + return 'video_failed: encoder or device connection failed; reconnect the selected phone and retry video' + } +} + +async function terminateChild(child: ChildProcess): Promise { + if (child.exitCode !== null) return + const exited = new Promise(resolve => child.once('exit', () => resolve())) + if (!child.killed) child.kill('SIGTERM') + const graceful = await Promise.race([ + exited.then(() => true), + new Promise(resolve => setTimeout(() => resolve(false), 500)), + ]) + if (!graceful && child.exitCode === null) { + child.kill('SIGKILL') + await Promise.race([exited, new Promise(resolve => setTimeout(resolve, 250))]) + } +} + +async function availableTcpPort(): Promise { + return new Promise((resolvePort, rejectPort) => { + const server = createServer() + server.once('error', rejectPort) + server.listen(0, '127.0.0.1', () => { + const address = server.address() + const port = typeof address === 'object' && address !== null ? address.port : 0 + server.close(error => error === undefined ? resolvePort(port) : rejectPort(error)) + }) + }) +} + +async function waitForSpawn(child: ChildProcess, signal: AbortSignal): Promise { + await new Promise((resolve, reject) => { + const done = (error?: Error): void => { + child.off('spawn', onSpawn) + child.off('error', onError) + signal.removeEventListener('abort', onAbort) + error === undefined ? resolve() : reject(error) + } + const onSpawn = (): void => done() + const onError = (error: Error): void => done(error) + const onAbort = (): void => done(signal.reason instanceof Error ? signal.reason : new Error(String(signal.reason))) + child.once('spawn', onSpawn) + child.once('error', onError) + signal.addEventListener('abort', onAbort, { once: true }) + }) +} + +async function connectVideo(connectImpl: typeof connect, port: number, signal: AbortSignal): Promise { + const deadline = Date.now() + 10_000 + while (true) { + signal.throwIfAborted() + try { + const socket = await new Promise((resolve, reject) => { + const socket = connectImpl({ host: '127.0.0.1', port }) + const cleanup = (): void => { + socket.off('connect', onConnect) + socket.off('error', onError) + signal.removeEventListener('abort', onAbort) + } + const onConnect = (): void => { cleanup(); resolve(socket) } + const onError = (error: Error): void => { cleanup(); socket.destroy(); reject(error) } + const onAbort = (): void => { cleanup(); socket.destroy(); reject(signal.reason) } + socket.once('connect', onConnect) + socket.once('error', onError) + signal.addEventListener('abort', onAbort, { once: true }) + }) + try { + await waitForVideoData(socket, signal, Math.max(1, deadline - Date.now())) + return socket + } catch (error) { + socket.destroy() + throw error + } + } catch (error) { + if (signal.aborted || Date.now() >= deadline) throw error + await new Promise(resolve => setTimeout(resolve, 120)) + } + } +} + +/** ADB forward accepts TCP before the device abstract socket exists, then closes it. + * Treat a connection as ready only after scrcpy has produced its first bytes. */ +async function waitForVideoData(socket: Socket, signal: AbortSignal, timeoutMs: number): Promise { + await new Promise((resolve, reject) => { + const timeout = setTimeout(() => done(new Error('opengui-workbuddy: timed out waiting for scrcpy video data')), timeoutMs) + const done = (error?: Error): void => { + clearTimeout(timeout) + socket.off('readable', onReadable) + socket.off('end', onClose) + socket.off('close', onClose) + socket.off('error', onError) + signal.removeEventListener('abort', onAbort) + error === undefined ? resolve() : reject(error) + } + const onReadable = (): void => { + if (socket.readableLength > 0) done() + } + const onClose = (): void => done(new Error('opengui-workbuddy: scrcpy video socket not ready')) + const onError = (error: Error): void => done(error) + const onAbort = (): void => done(signal.reason instanceof Error ? signal.reason : new Error(String(signal.reason))) + socket.once('readable', onReadable) + socket.once('end', onClose) + socket.once('close', onClose) + socket.once('error', onError) + signal.addEventListener('abort', onAbort, { once: true }) + }) +} diff --git a/packages/device-runtime/src/session-runtime.ts b/packages/device-runtime/src/session-runtime.ts new file mode 100644 index 0000000..f07ce6e --- /dev/null +++ b/packages/device-runtime/src/session-runtime.ts @@ -0,0 +1,80 @@ +/** Resource ownership is local to one host runtime, never machine-global. */ +export interface RuntimeDevice { readonly device: { readonly id: string; readonly serial: string } } +export interface RuntimeSession { + readonly id: string + state: string + readonly controller: AbortController + readonly pending: Set> + readonly devices: readonly RuntimeDevice[] +} + +/** Shared resource mechanics; adapters retain authentication and lease policy. */ +export class SessionRuntime { + readonly sessions = new Map() + readonly locks = new Map() + private readonly leaseOwners = new Map() + private readonly cleanups = new WeakMap>() + + register(record: R, control: boolean): void { + if (this.sessions.has(record.id)) throw new Error('opengui: duplicate sessionId') + const serials = record.devices.map(item => item.device.serial) + if (new Set(serials).size !== serials.length) throw new Error('opengui: a session cannot contain the same phone twice') + if (control && serials.some(serial => this.locks.has(serial))) throw new Error('opengui: device is already locked by another session') + this.sessions.set(record.id, record) + if (control) for (const serial of serials) { this.locks.set(serial, record.id); this.leaseOwners.set(serial, record) } + } + + requireSession(id: string): R { + const record = this.sessions.get(id) + if (!record) throw new Error('opengui: unknown sessionId') + return record + } + + requireActiveSession(id: string): R { + const record = this.requireSession(id) + if (record.state !== 'active') throw new Error(`opengui: session is ${record.state}`) + return record + } + + resolveDevice(record: R, deviceId: string | undefined): R['devices'][number] { + if (deviceId === undefined) { + if (record.devices.length !== 1) throw new Error('opengui: deviceId is required for a multi-device session') + return record.devices[0]! + } + const item = record.devices.find(candidate => candidate.device.id === deviceId) + if (!item) throw new Error('opengui: deviceId is not locked by this session') + return item + } + + release(record: R): void { + // A late callback cannot release a newly registered session with a reused id. + for (const item of record.devices) { + if (this.leaseOwners.get(item.device.serial) !== record) continue + this.locks.delete(item.device.serial) + this.leaseOwners.delete(item.device.serial) + } + } + + async track(record: R, operation: () => Promise): Promise { + const pending = Promise.resolve().then(() => { + record.controller.signal.throwIfAborted() + if (record.state !== 'active') throw new Error(`opengui: session is ${record.state}`) + return operation() + }) + record.pending.add(pending) + try { return await pending } finally { record.pending.delete(pending) } + } + + /** The caller terminates admission before draining; cleanup runs exactly once. */ + drain(record: R, releaseResources: () => Promise): Promise { + const existing = this.cleanups.get(record) + if (existing) return existing + if (record.state === 'active' && !record.controller.signal.aborted) throw new Error('opengui: stop session before cleanup') + const pending = (async () => { + await Promise.allSettled([...record.pending]) + try { await releaseResources() } finally { this.release(record) } + })() + this.cleanups.set(record, pending) + return pending + } +} diff --git a/packages/device-runtime/tests/session-contract.ts b/packages/device-runtime/tests/session-contract.ts new file mode 100644 index 0000000..c942c70 --- /dev/null +++ b/packages/device-runtime/tests/session-contract.ts @@ -0,0 +1,64 @@ +import assert from 'node:assert/strict' +import { SessionRuntime, type RuntimeSession } from '../src/session-runtime.ts' + +type Test = (name: string, body: () => Promise) => unknown +const session = (id: string, serial = 'serial-a'): RuntimeSession => ({ + id, state: 'active', controller: new AbortController(), pending: new Set(), + devices: [{ device: { id: 'phone-a', serial } }], +}) + +export function sessionContract(it: Test): void { + it('keeps device ownership until in-flight work and cleanup have drained', async () => { + const runtime = new SessionRuntime() + const a = session('a'); runtime.register(a, true) + let finish!: () => void + let entered!: () => void + const enteredPromise = new Promise(resolve => { entered = resolve }) + const task = runtime.track(a, async () => { entered(); await new Promise(resolve => { finish = resolve }) }) + await enteredPromise + a.state = 'closed'; a.controller.abort() + let releases = 0 + const cleanup = runtime.drain(a, async () => { releases++ }) + assert.equal(runtime.drain(a, async () => { releases++ }), cleanup) + assert.throws(() => runtime.register(session('b'), true), /locked/) + finish(); await task; await cleanup + assert.equal(releases, 1) + runtime.register(session('b'), true) + runtime.release(a) + assert.equal(runtime.locks.get('serial-a'), 'b') + }) + + it('does not release a new owner when a session id is reused', async () => { + const runtime = new SessionRuntime() + const old = session('same'); runtime.register(old, true) + runtime.release(old); runtime.sessions.delete(old.id) + const current = session('same'); runtime.register(current, true) + runtime.release(old) + assert.equal(runtime.locks.get('serial-a'), 'same') + runtime.release(current) + assert.equal(runtime.locks.size, 0) + }) + + it('rejects conflicting multi-device admission without claiming any device', async () => { + const runtime = new SessionRuntime() + runtime.register(session('a'), true) + const b = session('b', 'serial-b') + const multi = { ...b, devices: [...b.devices, ...session('other').devices] } + assert.throws(() => runtime.register(multi, true), /locked/) + assert.equal(runtime.sessions.has('b'), false) + assert.equal(runtime.locks.has('serial-b'), false) + assert.throws(() => runtime.register({ ...b, devices: [...b.devices, ...b.devices] }, true), /same phone/) + }) + + it('keeps independent runtime instances and prevents admission after cancellation', async () => { + const first = new SessionRuntime(), second = new SessionRuntime() + const a = session('a'), b = session('b') + first.register(a, true); second.register(b, true) + a.state = 'cancelled'; a.controller.abort(new Error('cancelled')) + let dispatched = false + await assert.rejects(first.track(a, async () => { dispatched = true }), /cancelled/) + await first.drain(a, async () => {}) + assert.equal(dispatched, false) + assert.equal(second.locks.get('serial-a'), 'b') + }) +} diff --git a/plugins/opengui/src/codex/service.ts b/plugins/opengui/src/codex/service.ts index eec6d1d..744f98d 100644 --- a/plugins/opengui/src/codex/service.ts +++ b/plugins/opengui/src/codex/service.ts @@ -1,3 +1,5 @@ +import { OpenGuiError } from '../errors.ts' +import { SessionRuntime } from '../../../../packages/device-runtime/src/session-runtime.ts' import { ViewerServer, type ViewerStreams } from '../viewer.ts' import { ScrcpyVideoStreams } from '../scrcpy-stream.ts' import { randomBytes, randomUUID } from 'node:crypto' @@ -276,8 +278,9 @@ export interface CodexOpenGuiServiceOptions { export class CodexOpenGuiService { private readonly host: CodexPhoneHost private readonly createSessionId: () => string - private readonly sessions = new Map() - private readonly locks = new Map() + private readonly runtime = new SessionRuntime() + private readonly sessions = this.runtime.sessions + private readonly locks = this.runtime.locks readonly viewers: ViewerServer private readonly wall: DeviceWallServer private readonly now: () => number @@ -343,9 +346,8 @@ export class CodexOpenGuiService { } for (const item of record.devices) { this.host.assignTarget(item.actor, item.device.serial) - if (mode === 'control') this.locks.set(item.device.serial, id) } - this.sessions.set(id, record) + this.runtime.register(record, mode === 'control') try { await this.wall.start() signal.throwIfAborted() @@ -381,7 +383,7 @@ export class CodexOpenGuiService { delete action.deviceId delete action.externalSideEffect delete action.confirmedExternalSideEffect - return this.runPhoneOperation(sessionId, deviceId, signal, (item, combined) => this.host.act(item.actor, action, combined)) + return this.runPhoneOperation(sessionId, deviceId, signal, (item, combined) => this.host.act(item.actor, action, combined), true) } async status(sessionId: string, signal: AbortSignal, renew = true): Promise { @@ -462,23 +464,24 @@ export class CodexOpenGuiService { deviceId: string | undefined, signal: AbortSignal, operation: (item: SessionDevice, combined: AbortSignal) => Promise, + mutating = false, ): Promise { const record = this.requireActiveSession(sessionId) this.viewers.assertReady(record.viewerId!) record.lastRequestAt = this.now() const item = this.resolveDevice(record, deviceId) const combined = AbortSignal.any([record.controller.signal, signal]) - const pending = operation(item, combined) - record.pending.add(pending) + const pending = this.runtime.track(record, () => operation(item, combined)) + let completed = false try { const value = await pending + completed = true combined.throwIfAborted() return this.publicObservation(record.id, item.device.id, value) } catch (error) { record.lastError = error instanceof Error ? error.message : String(error) + if (mutating && completed) throw new OpenGuiError('result_delivery_failed', record.lastError, 'outcome_unknown', 'observe') throw error - } finally { - record.pending.delete(pending) } } @@ -502,33 +505,15 @@ export class CodexOpenGuiService { } } - private requireSession(sessionId: string): SessionRecord { - const record = this.sessions.get(sessionId) - if (record === undefined) throw new Error('opengui: unknown sessionId') - return record - } + private requireSession(sessionId: string): SessionRecord { return this.runtime.requireSession(sessionId) } - private requireActiveSession(sessionId: string): SessionRecord { - const record = this.requireSession(sessionId) - if (record.state !== 'active') throw new Error(`opengui: session is ${record.state}`) - return record - } + private requireActiveSession(sessionId: string): SessionRecord { return this.runtime.requireActiveSession(sessionId) } private resolveDevice(record: SessionRecord, deviceId: string | undefined): SessionDevice { - if (deviceId === undefined) { - if (record.devices.length !== 1) throw new Error('opengui: deviceId is required for a multi-device session') - return record.devices[0]! - } - const item = record.devices.find(candidate => candidate.device.id === deviceId) - if (item === undefined) throw new Error('opengui: deviceId is not locked by this session') - return item + return this.runtime.resolveDevice(record, deviceId) } - private release(record: SessionRecord): void { - for (const item of record.devices) { - if (this.locks.get(item.device.serial) === record.id) this.locks.delete(item.device.serial) - } - } + private release(record: SessionRecord): void { this.runtime.release(record) } private async releaseDeviceResources(record: SessionRecord): Promise { if (record.mode === 'observe') return @@ -571,14 +556,12 @@ export class CodexOpenGuiService { record.state = state record.closedAt = new Date(this.now()).toISOString() record.controller.abort(new Error('opengui: session ' + state)) - record.finishing = (async () => { - // Keep the exclusive lease until old work and its cleanup have finished. - await Promise.allSettled(record.pending) + record.finishing = this.runtime.drain(record, async () => { await this.releaseDeviceResources(record) try { await this.onSessionClosed(record.id) } catch (error) { record.lastError = error instanceof Error ? error.message : String(error) } - finally { this.release(record); this.pruneClosedSessions() } - })() + finally { this.pruneClosedSessions() } + }) return record.finishing } } diff --git a/plugins/opengui/src/daemon.ts b/plugins/opengui/src/daemon.ts index 07a5d30..508e363 100644 --- a/plugins/opengui/src/daemon.ts +++ b/plugins/opengui/src/daemon.ts @@ -1,4 +1,4 @@ -import { errorInfo } from './errors.ts' +import { errorInfo, OpenGuiError } from './errors.ts' import { randomUUID } from 'node:crypto' import { spawn } from 'node:child_process' import { chmod, lstat, open, readFile, rm } from 'node:fs/promises' @@ -31,22 +31,28 @@ export function sendRequest(endpoint: string, value: Request, signal?: AbortSign return new Promise((resolve, reject) => { const socket = createConnection(endpoint) let body = '' + let sent = false + const fail = (error: unknown): void => { + cleanup() + reject(value.name === 'opengui_act' && sent + ? new OpenGuiError('connection_lost', error instanceof Error ? error.message : String(error), 'outcome_unknown', 'observe') : error) + } const cleanup = (): void => { signal?.removeEventListener('abort', abort); socket.destroy() } - const abort = (): void => { cleanup(); reject(signal?.reason ?? new Error('opengui: request cancelled')) } + const abort = (): void => { fail(signal?.reason ?? new Error('opengui: request cancelled')) } if (signal?.aborted) { abort(); return } signal?.addEventListener('abort', abort, { once: true }) socket.setEncoding('utf8') - socket.setTimeout(125_000, () => { cleanup(); reject(new Error('opengui: daemon request timed out')) }) - socket.once('connect', () => socket.write(JSON.stringify(value) + '\n')) + socket.setTimeout(125_000, () => fail(new Error('opengui: daemon request timed out'))) + socket.once('connect', () => { sent = true; socket.write(JSON.stringify(value) + '\n') }) socket.on('data', chunk => { body += chunk - if (Buffer.byteLength(body) > 2_000_000) { cleanup(); reject(new Error('opengui: oversized daemon response')); return } + if (Buffer.byteLength(body) > 2_000_000) { fail(new Error('opengui: oversized daemon response')); return } if (!body.includes('\n')) return try { const result = JSON.parse(body.trim()) as Response; cleanup(); resolve(result) } - catch (error) { cleanup(); reject(error) } + catch (error) { fail(error) } }) - socket.once('error', error => { cleanup(); reject(error) }) - socket.once('end', () => { if (!body.includes('\n')) { cleanup(); reject(new Error('opengui: incomplete daemon response')) } }) + socket.once('error', fail) + socket.once('end', () => { if (!body.includes('\n')) { fail(new Error('opengui: incomplete daemon response')) } }) }) } diff --git a/plugins/opengui/src/forward-registry.ts b/plugins/opengui/src/forward-registry.ts index d9713e0..2f34ac5 100644 --- a/plugins/opengui/src/forward-registry.ts +++ b/plugins/opengui/src/forward-registry.ts @@ -1,205 +1 @@ -import { randomUUID } from 'node:crypto' -import { mkdir, open, readFile, rename, rm, stat, writeFile } from 'node:fs/promises' -import { dirname } from 'node:path' - -export type OwnedForwardKind = 'text-input' | 'video-stream' - -export interface OwnedForward { - readonly serial: string - readonly port: number - readonly scid: string - readonly kind: OwnedForwardKind -} - -interface StoredForward extends OwnedForward { - readonly ownerId?: string - readonly ownerPid?: number -} - -export type ForwardAdbRunner = (args: readonly string[], signal: AbortSignal) => Promise - -interface ListedForward { - readonly serial: string - readonly local: string - readonly remote: string -} - -export function parseAdbForwardList(output: string): ListedForward[] { - const forwards: ListedForward[] = [] - for (const line of output.split(/\r?\n/u)) { - const [serial, local, remote, ...extra] = line.trim().split(/\s+/u) - if (!serial || !local || !remote || extra.length > 0) continue - forwards.push({ serial, local, remote }) - } - return forwards -} - - -/** Durable inventory of ADB forwards created by this plugin, never inferred from global ADB state. */ -export class OwnedForwardRegistry { - private mutation = Promise.resolve() - private readonly ownerId = randomUUID() - - constructor(private readonly path: string) {} - - async track(record: OwnedForward): Promise { - await this.mutate(records => { - records.set(this.key(record), { ...record, ownerId: this.ownerId, ownerPid: process.pid }) - }) - } - - async release(record: OwnedForward, runAdb: ForwardAdbRunner, timeoutMs = 5_000): Promise { - let forwards: ListedForward[] - try { - const listed = await runAdb( - ['-s', record.serial, 'forward', '--list'], - AbortSignal.timeout(timeoutMs), - ) - forwards = parseAdbForwardList(String(listed ?? '')) - } catch { - return false - } - const local = `tcp:${record.port}` - const owned = forwards.find(candidate => candidate.serial === record.serial && candidate.local === local) - if (owned === undefined || owned.remote !== this.remote(record)) { - try { - await this.deleteMatching(record) - return true - } catch { - return false - } - } - try { - await runAdb( - ['-s', record.serial, 'forward', '--remove', local], - AbortSignal.timeout(timeoutMs), - ) - } catch { - return false - } - try { - await this.deleteMatching(record) - return true - } catch { - return false - } - } - - async recover(runAdb: ForwardAdbRunner): Promise<{ removed: number; retained: number }> { - const records = await this.read() - let removed = 0 - let retained = 0 - for (const record of records.values()) { - if (record.ownerId !== undefined && record.ownerId !== this.ownerId - && record.ownerPid !== undefined && this.processAlive(record.ownerPid)) { - retained += 1 - continue - } - if (await this.release(record, runAdb)) removed += 1 - else retained += 1 - } - return { removed, retained } - } - - async list(): Promise { - return [...(await this.read()).values()].map(record => this.publicRecord(record)) - } - - private key(record: OwnedForward): string { - return `${record.serial}\u0000${record.port}` - } - - private remote(record: OwnedForward): string { - return `localabstract:scrcpy_${record.scid}` - } - - private valid(value: unknown): value is StoredForward { - if (typeof value !== 'object' || value === null) return false - const candidate = value as Partial - return typeof candidate.serial === 'string' && Number.isSafeInteger(candidate.port) - && typeof candidate.scid === 'string' - && (candidate.kind === 'text-input' || candidate.kind === 'video-stream') - } - - private async mutate(change: (records: Map) => void | boolean): Promise { - const operation = this.mutation.then(async () => { - await mkdir(dirname(this.path), { recursive: true }) - const releaseLock = await this.acquireLock() - try { - const records = await this.read() - change(records) - if (records.size === 0) { - await rm(this.path, { force: true }) - return - } - const temporary = `${this.path}.${process.pid}.${this.ownerId}.${randomUUID()}.tmp` - await writeFile(temporary, JSON.stringify([...records.values()]), { encoding: 'utf8', mode: 0o600 }) - await rename(temporary, this.path) - } finally { - await releaseLock() - } - }) - this.mutation = operation.catch(() => undefined) - await operation - } - - private async read(): Promise> { - try { - const parsed = JSON.parse(await readFile(this.path, 'utf8')) as unknown - if (!Array.isArray(parsed)) return new Map() - const records = new Map() - for (const value of parsed) { - if (!this.valid(value)) continue - const record = value - records.set(this.key(record), record) - } - return records - } catch { - return new Map() - } - } - - private async deleteMatching(record: OwnedForward): Promise { - await this.mutate(records => { - const stored = records.get(this.key(record)) - if (stored?.scid === record.scid) records.delete(this.key(record)) - }) - } - - private publicRecord(record: StoredForward): OwnedForward { - return { serial: record.serial, port: record.port, scid: record.scid, kind: record.kind } - } - - private processAlive(pid: number): boolean { - try { - process.kill(pid, 0) - return true - } catch { - return false - } - } - - private async acquireLock(): Promise<() => Promise> { - const lock = `${this.path}.lock` - const deadline = Date.now() + 5_000 - while (true) { - try { - const handle = await open(lock, 'wx', 0o600) - return async () => { - await handle.close() - await rm(lock, { force: true }) - } - } catch (error) { - if ((error as NodeJS.ErrnoException).code !== 'EEXIST') throw error - try { - if (Date.now() - (await stat(lock)).mtimeMs > 10_000) { - await rm(lock, { force: true }) - continue - } - } catch { /* another writer released the lock */ } - if (Date.now() >= deadline) throw new Error('opengui: timed out waiting for forward registry lock') - await new Promise(resolve => setTimeout(resolve, 10)) - } - } - } -} +export * from '../../../packages/device-runtime/src/forward-registry.ts' diff --git a/plugins/opengui/src/scrcpy-stream.ts b/plugins/opengui/src/scrcpy-stream.ts index 524a1df..3ddcd6f 100644 --- a/plugins/opengui/src/scrcpy-stream.ts +++ b/plugins/opengui/src/scrcpy-stream.ts @@ -1,543 +1,15 @@ -// Ported from the repository video transport; see VIDEO-NOTICE.md. -import { randomBytes } from 'node:crypto' -import { spawn } from 'node:child_process' -import type { ChildProcess } from 'node:child_process' -import { connect, createServer } from 'node:net' -import type { Socket } from 'node:net' -export interface VideoDevice { readonly id: string; readonly serial: string } -import { OwnedForwardRegistry } from './forward-registry.ts' -import type { OwnedForward } from './forward-registry.ts' -import { - SCRCPY_VERSION, - type ScrcpyAsset, - ScrcpyInstaller, - resolveScrcpyAsset, -} from './scrcpy.ts' - -const SCRCPY_REMOTE_SERVER = '/data/local/tmp/opengui-codex-scrcpy-server.jar' -const SESSION_PACKET_FLAG = 0x8000000000000000n -const CONFIG_PACKET_FLAG = 0x4000000000000000n -const KEY_PACKET_FLAG = 0x2000000000000000n -const PTS_MASK = 0x1fffffffffffffffn - -export type ScrcpyVideoEvent = { - readonly type: 'codec' - readonly codec: 'h264' -} | { - readonly type: 'session' - readonly width: number - readonly height: number - readonly clientResized: boolean -} | { - readonly type: 'packet' - readonly config: boolean - readonly key: boolean - readonly pts: bigint - readonly data: Buffer -} - -/** Incremental parser for scrcpy 4.1 stream metadata and H.264 media packets. */ -export class ScrcpyVideoPacketParser { - private buffer = Buffer.alloc(0) - private codecRead = false - - push(chunk: Buffer): ScrcpyVideoEvent[] { - if (chunk.length > 0) this.buffer = Buffer.concat([this.buffer, chunk]) - const events: ScrcpyVideoEvent[] = [] - if (!this.codecRead) { - if (this.buffer.length < 4) return events - const codec = this.buffer.subarray(0, 4).toString('ascii') - if (codec !== 'h264') throw new Error(`opengui-codex: unsupported scrcpy video codec ${codec}`) - this.codecRead = true - this.buffer = this.buffer.subarray(4) - events.push({ type: 'codec', codec: 'h264' }) - } - while (this.buffer.length >= 12) { - const flagsAndPts = this.buffer.readBigUInt64BE(0) - if ((flagsAndPts & SESSION_PACKET_FLAG) !== 0n) { - const flags = this.buffer.readUInt32BE(0) - const width = this.buffer.readUInt32BE(4) - const height = this.buffer.readUInt32BE(8) - if (width < 1 || height < 1 || width > 16_384 || height > 16_384) { - throw new Error(`opengui-codex: invalid scrcpy video size ${width}x${height}`) - } - this.buffer = this.buffer.subarray(12) - events.push({ type: 'session', width, height, clientResized: (flags & 1) === 1 }) - continue - } - const size = this.buffer.readUInt32BE(8) - if (size > 16 * 1024 * 1024) throw new Error('opengui-codex: scrcpy video packet exceeds 16 MiB') - if (this.buffer.length < 12 + size) break - const data = Buffer.from(this.buffer.subarray(12, 12 + size)) - this.buffer = this.buffer.subarray(12 + size) - events.push({ - type: 'packet', - config: (flagsAndPts & CONFIG_PACKET_FLAG) !== 0n, - key: (flagsAndPts & KEY_PACKET_FLAG) !== 0n, - pts: flagsAndPts & PTS_MASK, - data, - }) - } - return events - } -} - -/** Fixed read-only server options for the embedded low-latency stream. */ -export function buildScrcpyVideoServerArgs(scid: string, serverPath = SCRCPY_REMOTE_SERVER): string[] { - return [ - `CLASSPATH=${serverPath}`, - 'app_process', '/', 'com.genymobile.scrcpy.Server', SCRCPY_VERSION, - `scid=${scid}`, - 'tunnel_forward=true', - 'video=true', - 'audio=false', - 'control=false', - 'cleanup=false', - 'video_codec=h264', - 'max_size=960', - 'max_fps=30', - 'video_bit_rate=2000000', - 'video_codec_options=i-frame-interval=1', - 'send_dummy_byte=false', - 'send_device_meta=false', - 'send_stream_meta=true', - 'send_frame_meta=true', - ] -} - -export interface ScrcpyStreamSink { - sendText(text: string): void - sendBinary(data: Buffer): void - bufferedBytes(): number - close(code?: number, reason?: string): void - onClose(listener: () => void): void -} - -export interface ScrcpyStreamStatus { - supported: boolean - cached: boolean - /** @deprecated Kept true on supported Hosts for one client-compatibility release. */ - approved: boolean - phase: 'idle' | 'downloading' | 'extracting' | 'ready' | 'error' - version: string - totalBytes?: number - downloadedBytes?: number - activeSources: number - maxSources: number - message?: string -} - -type AdbRunner = (args: readonly string[], signal: AbortSignal) => Promise - -interface StreamEntry { - readonly device: VideoDevice - readonly subscribers: Set - readonly waiting: Set - readonly controller: AbortController - operation: Promise - socket?: Socket - process?: ChildProcess - port?: number - forward?: OwnedForward - idleTimer: ReturnType | undefined - lastCodec?: string - lastSession?: string - replay: Buffer[] - replayBytes: number - closeCode: number - closeReason: string - closed: boolean - closing?: Promise -} - -export interface ScrcpyVideoStreamsOptions { - adbPath: () => string - runAdb: AdbRunner - installer: ScrcpyInstaller - asset?: ScrcpyAsset - spawn?: typeof spawn - connect?: typeof connect - freePort?: () => Promise - idleGraceMs?: number - maxSources?: number - onError?: (error: unknown) => void - forwardRegistry: OwnedForwardRegistry -} - -/** Shares one scrcpy encoder per device across same-origin browser subscribers. */ -export class ScrcpyVideoStreams { - private readonly installer: ScrcpyInstaller - private readonly asset: ScrcpyAsset | undefined - private readonly spawnImpl: typeof spawn - private readonly connectImpl: typeof connect - private readonly freePort: () => Promise - private readonly idleGraceMs: number - private readonly maxSources: number - private readonly onError: (error: unknown) => void - private readonly forwardRegistry: OwnedForwardRegistry - private readonly entries = new Map() - private readonly lifetime = new AbortController() - private phase: ScrcpyStreamStatus['phase'] = 'idle' - private downloadedBytes: number | undefined - private message: string | undefined - - constructor(private readonly options: ScrcpyVideoStreamsOptions) { - this.installer = options.installer - this.asset = options.asset ?? resolveScrcpyAsset() - this.spawnImpl = options.spawn ?? spawn - this.connectImpl = options.connect ?? connect - this.freePort = options.freePort ?? availableTcpPort - this.idleGraceMs = options.idleGraceMs ?? 2_000 - this.maxSources = options.maxSources ?? 4 - this.onError = options.onError ?? (() => {}) - this.forwardRegistry = options.forwardRegistry - } - - async prepare(signal: AbortSignal): Promise { - if (!this.asset) throw new Error('stream_unsupported') - await this.installer.ensure(this.asset, signal, () => {}) - } - - approve(): boolean { - // Compatibility endpoint: first-use preparation is automatic now. - return this.asset !== undefined - } - - async status(): Promise { - const cached = this.asset !== undefined && await this.installer.isInstalled(this.asset) - if (cached && this.phase === 'idle') this.phase = 'ready' - return { - supported: this.asset !== undefined, - cached, - approved: this.asset !== undefined, - phase: cached && this.phase === 'idle' ? 'ready' : this.phase, - version: SCRCPY_VERSION, - ...(this.asset === undefined ? {} : { totalBytes: this.asset.bytes }), - ...(this.downloadedBytes === undefined ? {} : { downloadedBytes: this.downloadedBytes }), - activeSources: this.entries.size, - maxSources: this.maxSources, - ...(this.message === undefined ? {} : { message: this.message }), - } - } - - async subscribe(device: VideoDevice, sink: ScrcpyStreamSink): Promise<() => void> { - if (this.lifetime.signal.aborted) throw new Error('stream_disposed') - const asset = this.asset - if (asset === undefined) throw new Error('stream_unsupported') - let entry = this.entries.get(device.id) - if (entry?.closed === true) { - // A closing encoder still consumes a source slot until its owned resources drain. - await entry.closing - return this.subscribe(device, sink) - } - if (entry === undefined) { - if (this.entries.size >= this.maxSources) throw new Error('stream_capacity_wait') - const controller = new AbortController() - entry = { - device, - subscribers: new Set(), - waiting: new Set(), - controller, - operation: Promise.resolve(), - idleTimer: undefined, - replay: [], - replayBytes: 0, - closeCode: 1000, - closeReason: 'stream stopped', - closed: false, - } - this.entries.set(device.id, entry) - entry.operation = this.start(entry, asset).catch(error => { - if (!entry!.controller.signal.aborted) { - this.phase = 'error' - this.message = error instanceof Error ? error.message : String(error) - this.onError(error) - this.broadcastText(entry!, { type: 'error', message: this.publicError(error) }) - entry!.closeCode = 1011 - entry!.closeReason = 'stream failed' - } - }).finally(() => { - void this.closeEntry(entry!) - }) - } - if (entry.idleTimer !== undefined) { - clearTimeout(entry.idleTimer) - entry.idleTimer = undefined - } - entry.subscribers.add(sink) - if (entry.lastCodec !== undefined) sink.sendText(entry.lastCodec) - if (entry.lastSession !== undefined) sink.sendText(entry.lastSession) - if (entry.replayBytes <= 1_000_000) { - for (const frame of entry.replay) sink.sendBinary(frame) - } else entry.waiting.add(sink) - return () => this.unsubscribe(entry!, sink) - } - - async dispose(): Promise { - if (!this.lifetime.signal.aborted) this.lifetime.abort(new Error('opengui-codex: stream manager disposed')) - await Promise.allSettled([...this.entries.values()].map(entry => this.closeEntry(entry))) - } - - private unsubscribe(entry: StreamEntry, sink: ScrcpyStreamSink): void { - entry.subscribers.delete(sink) - entry.waiting.delete(sink) - if (entry.subscribers.size > 0 || entry.closed || entry.idleTimer !== undefined) return - entry.idleTimer = setTimeout(() => { - entry.idleTimer = undefined - if (entry.subscribers.size === 0) void this.closeEntry(entry) - }, this.idleGraceMs) +import { ScrcpyVideoStreams as RuntimeStreams, buildScrcpyVideoServerArgs as serverArgs, + type ScrcpyVideoStreamsOptions as RuntimeOptions } from '../../../packages/device-runtime/src/scrcpy-stream.ts' +import { SCRCPY_VERSION, resolveScrcpyAsset } from './scrcpy.ts' +export { ScrcpyVideoPacketParser } from '../../../packages/device-runtime/src/scrcpy-stream.ts' +export type { VideoDevice, ScrcpyStreamSink, ScrcpyStreamStatus, ScrcpyVideoEvent } from '../../../packages/device-runtime/src/scrcpy-stream.ts' +export type ScrcpyVideoStreamsOptions = Omit +const remoteServer = '/data/local/tmp/opengui-codex-scrcpy-server.jar' +export function buildScrcpyVideoServerArgs(scid: string, serverPath = remoteServer): string[] { + return serverArgs(scid, serverPath, SCRCPY_VERSION) +} +export class ScrcpyVideoStreams extends RuntimeStreams { + constructor(options: ScrcpyVideoStreamsOptions) { + super({ ...options, asset: options.asset ?? resolveScrcpyAsset(), version: SCRCPY_VERSION, remoteServer }) } - - private async start(entry: StreamEntry, asset: ScrcpyAsset): Promise { - const signal = AbortSignal.any([entry.controller.signal, this.lifetime.signal]) - this.phase = 'downloading' - this.message = undefined - const installed = await this.installer.ensure(asset, signal, progress => { - this.phase = progress.phase - this.downloadedBytes = progress.downloadedBytes - this.broadcastText(entry, { type: 'install', phase: progress.phase, downloadedBytes: progress.downloadedBytes, totalBytes: progress.totalBytes }) - }) - this.phase = 'ready' - this.downloadedBytes = undefined - signal.throwIfAborted() - - const port = await this.freePort() - entry.port = port - const scid = (randomBytes(4).readUInt32BE(0) & 0x7fffffff).toString(16).padStart(8, '0') - const forward: OwnedForward = { serial: entry.device.serial, port, scid, kind: 'video-stream' } - await this.options.runAdb(['-s', entry.device.serial, 'push', installed.server, SCRCPY_REMOTE_SERVER], signal) - try { - await this.forwardRegistry.track(forward) - entry.forward = forward - await this.options.runAdb([ - '-s', entry.device.serial, 'forward', '--no-rebind', `tcp:${port}`, `localabstract:scrcpy_${scid}`, - ], signal) - } catch (error) { - await this.forwardRegistry.release(forward, this.options.runAdb).catch(() => false) - throw error - } - - const child = this.spawnImpl(this.options.adbPath(), [ - '-s', entry.device.serial, 'shell', ...buildScrcpyVideoServerArgs(scid), - ], { shell: false, windowsHide: true, stdio: ['ignore', 'ignore', 'pipe'] }) - entry.process = child - let stderr = '' - child.stderr?.on('data', (chunk: Buffer | string) => { stderr = `${stderr}${String(chunk)}`.slice(-2_000) }) - await waitForSpawn(child, signal) - const socket = await connectVideo(this.connectImpl, port, signal) - entry.socket = socket - const parser = new ScrcpyVideoPacketParser() - socket.on('data', chunk => { - try { - for (const event of parser.push(Buffer.from(chunk))) this.broadcastEvent(entry, event) - } catch (error) { - this.broadcastText(entry, { type: 'error', message: this.publicError(error) }) - void this.closeEntry(entry) - } - }) - const settled = new Promise((resolve, reject) => { - let socketCloseTimer: ReturnType | undefined - const rejectSocketClose = (): void => { - socketCloseTimer = setTimeout(() => { - reject(new Error(`scrcpy video socket closed unexpectedly${stderr.trim() ? `: ${stderr.trim()}` : ''}`)) - }, 150) - } - socket.once('error', reject) - socket.once('close', () => signal.aborted ? resolve() : rejectSocketClose()) - child.once('exit', (code, exitSignal) => { - if (socketCloseTimer !== undefined) clearTimeout(socketCloseTimer) - if (entry.controller.signal.aborted) resolve() - else reject(new Error(`scrcpy video server exited (code=${String(code)}, signal=${String(exitSignal)})${stderr.trim() ? `: ${stderr.trim()}` : ''}`)) - }) - signal.addEventListener('abort', () => { - if (socketCloseTimer !== undefined) clearTimeout(socketCloseTimer) - resolve() - }, { once: true }) - }) - await settled - } - - private broadcastEvent(entry: StreamEntry, event: ScrcpyVideoEvent): void { - if (event.type === 'packet') { - const frame = this.packetFrame(event) - if (event.config) { - entry.replay = [frame] - entry.replayBytes = frame.byteLength - } else if (event.key) { - entry.replay = [...entry.replay.filter(packet => (packet[0]! & 1) !== 0), frame] - entry.replayBytes = entry.replay.reduce((total, packet) => total + packet.byteLength, 0) - } else if (entry.replay.some(packet => (packet[0]! & 2) !== 0)) { - if (entry.replayBytes + frame.byteLength <= 8 * 1024 * 1024) { - entry.replay.push(frame) - entry.replayBytes += frame.byteLength - } else { - entry.replay = entry.replay.filter(packet => (packet[0]! & 1) !== 0) - entry.replayBytes = entry.replay.reduce((total, packet) => total + packet.byteLength, 0) - } - } - for (const sink of entry.subscribers) { - if (sink.bufferedBytes() > 1_000_000) { - entry.waiting.add(sink) - if (sink.bufferedBytes() > 2_000_000) sink.close(1013, 'slow_client') - continue - } - if (entry.waiting.has(sink)) { - if (!event.key) continue - sink.sendText(JSON.stringify({ type: 'reset' })) - for (const config of entry.replay.filter(packet => (packet[0]! & 1) !== 0)) sink.sendBinary(config) - entry.waiting.delete(sink) - } - sink.sendBinary(frame) - } - return - } - const text = JSON.stringify(event) - if (event.type === 'codec') entry.lastCodec = text - else { - entry.lastSession = text - entry.replay = [] - entry.replayBytes = 0 - } - for (const sink of entry.subscribers) sink.sendText(text) - } - - private packetFrame(event: Extract): Buffer { - const frame = Buffer.allocUnsafe(9 + event.data.length) - frame[0] = (event.config ? 1 : 0) | (event.key ? 2 : 0) - frame.writeBigUInt64BE(event.pts, 1) - event.data.copy(frame, 9) - return frame - } - - private broadcastText(entry: StreamEntry, value: unknown): void { - const text = JSON.stringify(value) - for (const sink of entry.subscribers) sink.sendText(text) - } - - private closeEntry(entry: StreamEntry): Promise { - if (entry.closing !== undefined) return entry.closing - if (entry.closed) return Promise.resolve() - entry.closed = true - const closing = (async (): Promise => { - if (entry.idleTimer !== undefined) clearTimeout(entry.idleTimer) - entry.controller.abort(new Error('opengui-codex: embedded stream stopped')) - entry.socket?.destroy() - const child = entry.process - if (child !== undefined && child.exitCode === null) await terminateChild(child) - if (entry.forward !== undefined) await this.forwardRegistry.release(entry.forward, this.options.runAdb, 1_000) - if (this.entries.get(entry.device.id) === entry) this.entries.delete(entry.device.id) - for (const sink of entry.subscribers) sink.close(entry.closeCode, entry.closeReason) - entry.subscribers.clear() - })() - entry.closing = closing - return closing - } - - private publicError(_error: unknown): string { - return 'video_failed: encoder or device connection failed; reconnect the selected phone and retry video' - } -} - -async function terminateChild(child: ChildProcess): Promise { - if (child.exitCode !== null) return - const exited = new Promise(resolve => child.once('exit', () => resolve())) - if (!child.killed) child.kill('SIGTERM') - const graceful = await Promise.race([ - exited.then(() => true), - new Promise(resolve => setTimeout(() => resolve(false), 500)), - ]) - if (!graceful && child.exitCode === null) { - child.kill('SIGKILL') - await Promise.race([exited, new Promise(resolve => setTimeout(resolve, 250))]) - } -} - -async function availableTcpPort(): Promise { - return new Promise((resolvePort, rejectPort) => { - const server = createServer() - server.once('error', rejectPort) - server.listen(0, '127.0.0.1', () => { - const address = server.address() - const port = typeof address === 'object' && address !== null ? address.port : 0 - server.close(error => error === undefined ? resolvePort(port) : rejectPort(error)) - }) - }) -} - -async function waitForSpawn(child: ChildProcess, signal: AbortSignal): Promise { - await new Promise((resolve, reject) => { - const done = (error?: Error): void => { - child.off('spawn', onSpawn) - child.off('error', onError) - signal.removeEventListener('abort', onAbort) - error === undefined ? resolve() : reject(error) - } - const onSpawn = (): void => done() - const onError = (error: Error): void => done(error) - const onAbort = (): void => done(signal.reason instanceof Error ? signal.reason : new Error(String(signal.reason))) - child.once('spawn', onSpawn) - child.once('error', onError) - signal.addEventListener('abort', onAbort, { once: true }) - }) -} - -async function connectVideo(connectImpl: typeof connect, port: number, signal: AbortSignal): Promise { - const deadline = Date.now() + 10_000 - while (true) { - signal.throwIfAborted() - try { - const socket = await new Promise((resolve, reject) => { - const socket = connectImpl({ host: '127.0.0.1', port }) - const cleanup = (): void => { - socket.off('connect', onConnect) - socket.off('error', onError) - signal.removeEventListener('abort', onAbort) - } - const onConnect = (): void => { cleanup(); resolve(socket) } - const onError = (error: Error): void => { cleanup(); socket.destroy(); reject(error) } - const onAbort = (): void => { cleanup(); socket.destroy(); reject(signal.reason) } - socket.once('connect', onConnect) - socket.once('error', onError) - signal.addEventListener('abort', onAbort, { once: true }) - }) - try { - await waitForVideoData(socket, signal, Math.max(1, deadline - Date.now())) - return socket - } catch (error) { - socket.destroy() - throw error - } - } catch (error) { - if (signal.aborted || Date.now() >= deadline) throw error - await new Promise(resolve => setTimeout(resolve, 120)) - } - } -} - -/** ADB forward accepts TCP before the device abstract socket exists, then closes it. - * Treat a connection as ready only after scrcpy has produced its first bytes. */ -async function waitForVideoData(socket: Socket, signal: AbortSignal, timeoutMs: number): Promise { - await new Promise((resolve, reject) => { - const timeout = setTimeout(() => done(new Error('opengui-codex: timed out waiting for scrcpy video data')), timeoutMs) - const done = (error?: Error): void => { - clearTimeout(timeout) - socket.off('readable', onReadable) - socket.off('end', onClose) - socket.off('close', onClose) - socket.off('error', onError) - signal.removeEventListener('abort', onAbort) - error === undefined ? resolve() : reject(error) - } - const onReadable = (): void => { - if (socket.readableLength > 0) done() - } - const onClose = (): void => done(new Error('opengui-codex: scrcpy video socket not ready')) - const onError = (error: Error): void => done(error) - const onAbort = (): void => done(signal.reason instanceof Error ? signal.reason : new Error(String(signal.reason))) - socket.once('readable', onReadable) - socket.once('end', onClose) - socket.once('close', onClose) - socket.once('error', onError) - signal.addEventListener('abort', onAbort, { once: true }) - }) } diff --git a/plugins/opengui/tests/daemon.spec.ts b/plugins/opengui/tests/daemon.spec.ts index 5d74e15..7c38c27 100644 --- a/plugins/opengui/tests/daemon.spec.ts +++ b/plugins/opengui/tests/daemon.spec.ts @@ -7,6 +7,7 @@ import { createConnection } from 'node:net' import { CodexOpenGuiService } from '../src/codex/service.ts' import { assertVersion, request as makeRequest, sendRequest, startDaemon } from '../src/daemon.ts' import { OpenGuiError } from '../src/errors.ts' +import { ObservationStore } from '../src/state.ts' import { FakeHost } from './fixtures.ts' const cleanup: (() => Promise)[] = [] @@ -28,6 +29,17 @@ async function open(endpoint: string): Promise { } describe('standalone daemon transport', () => { + it('does not report an action as unexecuted when saving its image fails', async () => { + const server = await daemon(), sessionId = await open(server.endpoint) + const save = vi.spyOn(ObservationStore.prototype, 'save').mockRejectedValueOnce(new Error('test disk full')) + try { + const response = await sendRequest(server.endpoint, request('opengui_act', { + sessionId, action: 'key', key: 'Home', observationId: 'frame', externalSideEffect: 'none', + })) + expect(response).toMatchObject({ ok: false, error: 'test disk full', failure: { executionState: 'outcome_unknown', recovery: 'observe' } }) + } finally { save.mockRestore() } + }) + it('preserves an unknown action outcome across the daemon transport', async () => { const server = await daemon(), sessionId = await open(server.endpoint) server.host.act = async () => { throw new OpenGuiError('capture_failed', 'capture failed after dispatch', 'outcome_unknown', 'observe') } diff --git a/plugins/opengui/tests/runtime-contract.spec.ts b/plugins/opengui/tests/runtime-contract.spec.ts index 7c255c6..cc50720 100644 --- a/plugins/opengui/tests/runtime-contract.spec.ts +++ b/plugins/opengui/tests/runtime-contract.spec.ts @@ -2,3 +2,5 @@ import { it } from 'vitest' import { PhoneController } from '../src/phone-controller.ts' import { controllerContract } from '../../../packages/device-runtime/tests/controller-contract.ts' controllerContract(it, options => new PhoneController(options)) +import { sessionContract } from '../../../packages/device-runtime/tests/session-contract.ts' +sessionContract(it) diff --git a/workbuddy-plugin/src/forward-registry.ts b/workbuddy-plugin/src/forward-registry.ts index 57d275f..d34a9fd 100644 --- a/workbuddy-plugin/src/forward-registry.ts +++ b/workbuddy-plugin/src/forward-registry.ts @@ -1,285 +1,8 @@ -import { randomUUID } from 'node:crypto' -import { mkdir, open, readFile, rename, rm, stat, writeFile } from 'node:fs/promises' -import { closeSync, openSync, readFileSync, renameSync, rmSync, writeFileSync } from 'node:fs' -import { spawnSync } from 'node:child_process' +import { join } from 'node:path' import { workbuddyStateDir } from './state.ts' -import { dirname, join } from 'node:path' - -export type OwnedForwardKind = 'text-input' | 'video-stream' - -export interface OwnedForward { - readonly serial: string - readonly port: number - readonly scid: string - readonly kind: OwnedForwardKind -} - -interface StoredForward extends OwnedForward { - readonly ownerId?: string - readonly ownerPid?: number -} - -export type ForwardAdbRunner = (args: readonly string[], signal: AbortSignal) => Promise - -interface ListedForward { - readonly serial: string - readonly local: string - readonly remote: string -} - -export function parseAdbForwardList(output: string): ListedForward[] { - const forwards: ListedForward[] = [] - for (const line of output.split(/\r?\n/u)) { - const [serial, local, remote, ...extra] = line.trim().split(/\s+/u) - if (!serial || !local || !remote || extra.length > 0) continue - forwards.push({ serial, local, remote }) - } - return forwards -} - -export function defaultForwardRegistryPath(): string { - return join(workbuddyStateDir(), 'owned-forwards.json') -} - -/** Durable inventory of ADB forwards created by this plugin, never inferred from global ADB state. */ -export class OwnedForwardRegistry { - private mutation = Promise.resolve() - private readonly ownerId = randomUUID() - - constructor(private readonly path = defaultForwardRegistryPath()) {} - - async track(record: OwnedForward): Promise { - await this.mutate(records => { - records.set(this.key(record), { ...record, ownerId: this.ownerId, ownerPid: process.pid }) - }) - } - - async release(record: OwnedForward, runAdb: ForwardAdbRunner, timeoutMs = 5_000): Promise { - let forwards: ListedForward[] - try { - const listed = await runAdb( - ['-s', record.serial, 'forward', '--list'], - AbortSignal.timeout(timeoutMs), - ) - forwards = parseAdbForwardList(String(listed ?? '')) - } catch { - return false - } - const local = `tcp:${record.port}` - const owned = forwards.find(candidate => candidate.serial === record.serial && candidate.local === local) - if (owned === undefined || owned.remote !== this.remote(record)) { - try { - await this.deleteMatching(record) - return true - } catch { - return false - } - } - try { - await runAdb( - ['-s', record.serial, 'forward', '--remove', local], - AbortSignal.timeout(timeoutMs), - ) - } catch { - return false - } - try { - await this.deleteMatching(record) - return true - } catch { - return false - } - } - - async recover(runAdb: ForwardAdbRunner): Promise<{ removed: number; retained: number }> { - const records = await this.read() - let removed = 0 - let retained = 0 - for (const record of records.values()) { - if (record.ownerId !== undefined && record.ownerId !== this.ownerId - && record.ownerPid !== undefined && this.processAlive(record.ownerPid)) { - retained += 1 - continue - } - if (await this.release(record, runAdb)) removed += 1 - else retained += 1 - } - return { removed, retained } - } - - async list(): Promise { - return [...(await this.read()).values()].map(record => this.publicRecord(record)) - } - - /** Best-effort synchronous signal-path cleanup when the Host cannot await plugin disposal. */ - releaseAllSync(adbPath: string): { removed: number; retained: number } { - const lock = `${this.path}.lock` - let lockFd: number - try { - lockFd = openSync(lock, 'wx', 0o600) - } catch { - return { removed: 0, retained: this.readSync().filter(record => record.ownerId === this.ownerId).length } - } - let records: StoredForward[] - try { - const value = JSON.parse(readFileSync(this.path, 'utf8')) as unknown - records = Array.isArray(value) ? value.filter(candidate => this.valid(candidate)) : [] - } catch { - closeSync(lockFd) - rmSync(lock, { force: true }) - return { removed: 0, retained: 0 } - } - const ownedRecords = records.filter(record => record.ownerId === this.ownerId) - const retainedOwned: StoredForward[] = [] - for (const record of ownedRecords) { - const listed = spawnSync(adbPath, ['-s', record.serial, 'forward', '--list'], { - timeout: 1_500, - encoding: 'utf8', - stdio: ['ignore', 'pipe', 'ignore'], - }) - if (listed.status !== 0 || listed.error !== undefined) { - retainedOwned.push(record) - continue - } - const local = `tcp:${record.port}` - const owned = parseAdbForwardList(listed.stdout ?? '') - .find(candidate => candidate.serial === record.serial && candidate.local === local) - if (owned === undefined || owned.remote !== this.remote(record)) continue - const removed = spawnSync(adbPath, ['-s', record.serial, 'forward', '--remove', local], { - timeout: 1_500, - stdio: 'ignore', - }) - if (removed.status !== 0 || removed.error !== undefined) retainedOwned.push(record) - } - const retainedKeys = new Set(retainedOwned.map(record => `${this.key(record)}\u0000${record.scid}`)) - const retained = records.filter(record => record.ownerId !== this.ownerId - || retainedKeys.has(`${this.key(record)}\u0000${record.scid}`)) - try { - if (retained.length === 0) rmSync(this.path, { force: true }) - else { - const temporary = `${this.path}.${process.pid}.${this.ownerId}.signal.tmp` - writeFileSync(temporary, JSON.stringify(retained), { encoding: 'utf8', mode: 0o600 }) - renameSync(temporary, this.path) - } - } catch { /* startup recovery retains the durable fallback */ } - finally { - closeSync(lockFd) - rmSync(lock, { force: true }) - } - return { removed: ownedRecords.length - retainedOwned.length, retained: retainedOwned.length } - } - - private key(record: OwnedForward): string { - return `${record.serial}\u0000${record.port}` - } - - private remote(record: OwnedForward): string { - return `localabstract:scrcpy_${record.scid}` - } - - private valid(value: unknown): value is StoredForward { - if (typeof value !== 'object' || value === null) return false - const candidate = value as Partial - return typeof candidate.serial === 'string' && Number.isSafeInteger(candidate.port) - && typeof candidate.scid === 'string' - && (candidate.kind === 'text-input' || candidate.kind === 'video-stream') - } - - private async mutate(change: (records: Map) => void | boolean): Promise { - const operation = this.mutation.then(async () => { - await mkdir(dirname(this.path), { recursive: true }) - const releaseLock = await this.acquireLock() - try { - const records = await this.read() - change(records) - if (records.size === 0) { - await rm(this.path, { force: true }) - return - } - const temporary = `${this.path}.${process.pid}.${this.ownerId}.${randomUUID()}.tmp` - await writeFile(temporary, JSON.stringify([...records.values()]), { encoding: 'utf8', mode: 0o600 }) - await rename(temporary, this.path) - } finally { - await releaseLock() - } - }) - this.mutation = operation.catch(() => undefined) - await operation - } - - private async read(): Promise> { - try { - const parsed = JSON.parse(await readFile(this.path, 'utf8')) as unknown - if (!Array.isArray(parsed)) return new Map() - const records = new Map() - for (const value of parsed) { - if (!this.valid(value)) continue - const record = value - records.set(this.key(record), record) - } - return records - } catch { - return new Map() - } - } - - private async deleteMatching(record: OwnedForward): Promise { - await this.mutate(records => { - const stored = records.get(this.key(record)) - if (stored?.scid === record.scid) records.delete(this.key(record)) - }) - } - - private publicRecord(record: StoredForward): OwnedForward { - return { serial: record.serial, port: record.port, scid: record.scid, kind: record.kind } - } - - private processAlive(pid: number): boolean { - try { - process.kill(pid, 0) - return true - } catch { - return false - } - } - - private readSync(): StoredForward[] { - try { - const parsed = JSON.parse(readFileSync(this.path, 'utf8')) as unknown - return Array.isArray(parsed) ? parsed.filter(candidate => this.valid(candidate)) : [] - } catch { - return [] - } - } - - private async acquireLock(): Promise<() => Promise> { - const lock = `${this.path}.lock` - const deadline = Date.now() + 5_000 - while (true) { - try { - const handle = await open(lock, 'wx', 0o600) - return async () => { - await handle.close() - await rm(lock, { force: true }) - } - } catch (error) { - // Windows can report a sharing violation while another writer deletes its lock. - // Retry acquisition only; never infer that EPERM grants permission to remove it. - if ((error as NodeJS.ErrnoException).code === 'EPERM') { - if (Date.now() >= deadline) throw error - await new Promise(resolve => setTimeout(resolve, 10)) - continue - } - if ((error as NodeJS.ErrnoException).code !== 'EEXIST') throw error - try { - if (Date.now() - (await stat(lock)).mtimeMs > 10_000) { - await rm(lock, { force: true }) - continue - } - } catch { /* another writer released the lock */ } - if (Date.now() >= deadline) throw new Error('opengui: timed out waiting for forward registry lock') - await new Promise(resolve => setTimeout(resolve, 10)) - } - } - } +import { OwnedForwardRegistry as RuntimeForwardRegistry } from '../../packages/device-runtime/src/forward-registry.ts' +export * from '../../packages/device-runtime/src/forward-registry.ts' +export function defaultForwardRegistryPath(): string { return join(workbuddyStateDir(), 'owned-forwards.json') } +export class OwnedForwardRegistry extends RuntimeForwardRegistry { + constructor(path = defaultForwardRegistryPath()) { super(path) } } diff --git a/workbuddy-plugin/src/scrcpy-stream.ts b/workbuddy-plugin/src/scrcpy-stream.ts index c7aab59..68ae852 100644 --- a/workbuddy-plugin/src/scrcpy-stream.ts +++ b/workbuddy-plugin/src/scrcpy-stream.ts @@ -1,543 +1,15 @@ -// Ported from the repository video transport; see VIDEO-NOTICE.md. -import { randomBytes } from 'node:crypto' -import { spawn } from 'node:child_process' -import type { ChildProcess } from 'node:child_process' -import { connect, createServer } from 'node:net' -import type { Socket } from 'node:net' -export interface VideoDevice { readonly id: string; readonly serial: string } -import { OwnedForwardRegistry } from './forward-registry.ts' -import type { OwnedForward } from './forward-registry.ts' -import { - SCRCPY_VERSION, - type ScrcpyAsset, - ScrcpyInstaller, - resolveScrcpyAsset, -} from './scrcpy.ts' - -const SCRCPY_REMOTE_SERVER = '/data/local/tmp/opengui-workbuddy-scrcpy-server.jar' -const SESSION_PACKET_FLAG = 0x8000000000000000n -const CONFIG_PACKET_FLAG = 0x4000000000000000n -const KEY_PACKET_FLAG = 0x2000000000000000n -const PTS_MASK = 0x1fffffffffffffffn - -export type ScrcpyVideoEvent = { - readonly type: 'codec' - readonly codec: 'h264' -} | { - readonly type: 'session' - readonly width: number - readonly height: number - readonly clientResized: boolean -} | { - readonly type: 'packet' - readonly config: boolean - readonly key: boolean - readonly pts: bigint - readonly data: Buffer -} - -/** Incremental parser for scrcpy 4.1 stream metadata and H.264 media packets. */ -export class ScrcpyVideoPacketParser { - private buffer = Buffer.alloc(0) - private codecRead = false - - push(chunk: Buffer): ScrcpyVideoEvent[] { - if (chunk.length > 0) this.buffer = Buffer.concat([this.buffer, chunk]) - const events: ScrcpyVideoEvent[] = [] - if (!this.codecRead) { - if (this.buffer.length < 4) return events - const codec = this.buffer.subarray(0, 4).toString('ascii') - if (codec !== 'h264') throw new Error(`opengui-workbuddy: unsupported scrcpy video codec ${codec}`) - this.codecRead = true - this.buffer = this.buffer.subarray(4) - events.push({ type: 'codec', codec: 'h264' }) - } - while (this.buffer.length >= 12) { - const flagsAndPts = this.buffer.readBigUInt64BE(0) - if ((flagsAndPts & SESSION_PACKET_FLAG) !== 0n) { - const flags = this.buffer.readUInt32BE(0) - const width = this.buffer.readUInt32BE(4) - const height = this.buffer.readUInt32BE(8) - if (width < 1 || height < 1 || width > 16_384 || height > 16_384) { - throw new Error(`opengui-workbuddy: invalid scrcpy video size ${width}x${height}`) - } - this.buffer = this.buffer.subarray(12) - events.push({ type: 'session', width, height, clientResized: (flags & 1) === 1 }) - continue - } - const size = this.buffer.readUInt32BE(8) - if (size > 16 * 1024 * 1024) throw new Error('opengui-workbuddy: scrcpy video packet exceeds 16 MiB') - if (this.buffer.length < 12 + size) break - const data = Buffer.from(this.buffer.subarray(12, 12 + size)) - this.buffer = this.buffer.subarray(12 + size) - events.push({ - type: 'packet', - config: (flagsAndPts & CONFIG_PACKET_FLAG) !== 0n, - key: (flagsAndPts & KEY_PACKET_FLAG) !== 0n, - pts: flagsAndPts & PTS_MASK, - data, - }) - } - return events - } -} - -/** Fixed read-only server options for the embedded low-latency stream. */ -export function buildScrcpyVideoServerArgs(scid: string, serverPath = SCRCPY_REMOTE_SERVER): string[] { - return [ - `CLASSPATH=${serverPath}`, - 'app_process', '/', 'com.genymobile.scrcpy.Server', SCRCPY_VERSION, - `scid=${scid}`, - 'tunnel_forward=true', - 'video=true', - 'audio=false', - 'control=false', - 'cleanup=false', - 'video_codec=h264', - 'max_size=960', - 'max_fps=30', - 'video_bit_rate=2000000', - 'video_codec_options=i-frame-interval=1', - 'send_dummy_byte=false', - 'send_device_meta=false', - 'send_stream_meta=true', - 'send_frame_meta=true', - ] -} - -export interface ScrcpyStreamSink { - sendText(text: string): void - sendBinary(data: Buffer): void - bufferedBytes(): number - close(code?: number, reason?: string): void - onClose(listener: () => void): void -} - -export interface ScrcpyStreamStatus { - supported: boolean - cached: boolean - /** @deprecated Kept true on supported Hosts for one client-compatibility release. */ - approved: boolean - phase: 'idle' | 'downloading' | 'extracting' | 'ready' | 'error' - version: string - totalBytes?: number - downloadedBytes?: number - activeSources: number - maxSources: number - message?: string -} - -type AdbRunner = (args: readonly string[], signal: AbortSignal) => Promise - -interface StreamEntry { - readonly device: VideoDevice - readonly subscribers: Set - readonly waiting: Set - readonly controller: AbortController - operation: Promise - socket?: Socket - process?: ChildProcess - port?: number - forward?: OwnedForward - idleTimer: ReturnType | undefined - lastCodec?: string - lastSession?: string - replay: Buffer[] - replayBytes: number - closeCode: number - closeReason: string - closed: boolean - closing?: Promise -} - -export interface ScrcpyVideoStreamsOptions { - adbPath: () => string - runAdb: AdbRunner - installer: ScrcpyInstaller - asset?: ScrcpyAsset - spawn?: typeof spawn - connect?: typeof connect - freePort?: () => Promise - idleGraceMs?: number - maxSources?: number - onError?: (error: unknown) => void - forwardRegistry: OwnedForwardRegistry -} - -/** Shares one scrcpy encoder per device across same-origin browser subscribers. */ -export class ScrcpyVideoStreams { - private readonly installer: ScrcpyInstaller - private readonly asset: ScrcpyAsset | undefined - private readonly spawnImpl: typeof spawn - private readonly connectImpl: typeof connect - private readonly freePort: () => Promise - private readonly idleGraceMs: number - private readonly maxSources: number - private readonly onError: (error: unknown) => void - private readonly forwardRegistry: OwnedForwardRegistry - private readonly entries = new Map() - private readonly lifetime = new AbortController() - private phase: ScrcpyStreamStatus['phase'] = 'idle' - private downloadedBytes: number | undefined - private message: string | undefined - - constructor(private readonly options: ScrcpyVideoStreamsOptions) { - this.installer = options.installer - this.asset = options.asset ?? resolveScrcpyAsset() - this.spawnImpl = options.spawn ?? spawn - this.connectImpl = options.connect ?? connect - this.freePort = options.freePort ?? availableTcpPort - this.idleGraceMs = options.idleGraceMs ?? 2_000 - this.maxSources = options.maxSources ?? 4 - this.onError = options.onError ?? (() => {}) - this.forwardRegistry = options.forwardRegistry - } - - async prepare(signal: AbortSignal): Promise { - if (!this.asset) throw new Error('stream_unsupported') - await this.installer.ensure(this.asset, signal, () => {}) - } - - approve(): boolean { - // Compatibility endpoint: first-use preparation is automatic now. - return this.asset !== undefined - } - - async status(): Promise { - const cached = this.asset !== undefined && await this.installer.isInstalled(this.asset) - if (cached && this.phase === 'idle') this.phase = 'ready' - return { - supported: this.asset !== undefined, - cached, - approved: this.asset !== undefined, - phase: cached && this.phase === 'idle' ? 'ready' : this.phase, - version: SCRCPY_VERSION, - ...(this.asset === undefined ? {} : { totalBytes: this.asset.bytes }), - ...(this.downloadedBytes === undefined ? {} : { downloadedBytes: this.downloadedBytes }), - activeSources: this.entries.size, - maxSources: this.maxSources, - ...(this.message === undefined ? {} : { message: this.message }), - } - } - - async subscribe(device: VideoDevice, sink: ScrcpyStreamSink): Promise<() => void> { - if (this.lifetime.signal.aborted) throw new Error('stream_disposed') - const asset = this.asset - if (asset === undefined) throw new Error('stream_unsupported') - let entry = this.entries.get(device.id) - if (entry?.closed === true) { - // A closing encoder still consumes a source slot until its owned resources drain. - await entry.closing - return this.subscribe(device, sink) - } - if (entry === undefined) { - if (this.entries.size >= this.maxSources) throw new Error('stream_capacity_wait') - const controller = new AbortController() - entry = { - device, - subscribers: new Set(), - waiting: new Set(), - controller, - operation: Promise.resolve(), - idleTimer: undefined, - replay: [], - replayBytes: 0, - closeCode: 1000, - closeReason: 'stream stopped', - closed: false, - } - this.entries.set(device.id, entry) - entry.operation = this.start(entry, asset).catch(error => { - if (!entry!.controller.signal.aborted) { - this.phase = 'error' - this.message = error instanceof Error ? error.message : String(error) - this.onError(error) - this.broadcastText(entry!, { type: 'error', message: this.publicError(error) }) - entry!.closeCode = 1011 - entry!.closeReason = 'stream failed' - } - }).finally(() => { - void this.closeEntry(entry!) - }) - } - if (entry.idleTimer !== undefined) { - clearTimeout(entry.idleTimer) - entry.idleTimer = undefined - } - entry.subscribers.add(sink) - if (entry.lastCodec !== undefined) sink.sendText(entry.lastCodec) - if (entry.lastSession !== undefined) sink.sendText(entry.lastSession) - if (entry.replayBytes <= 1_000_000) { - for (const frame of entry.replay) sink.sendBinary(frame) - } else entry.waiting.add(sink) - return () => this.unsubscribe(entry!, sink) - } - - async dispose(): Promise { - if (!this.lifetime.signal.aborted) this.lifetime.abort(new Error('opengui-workbuddy: stream manager disposed')) - await Promise.allSettled([...this.entries.values()].map(entry => this.closeEntry(entry))) - } - - private unsubscribe(entry: StreamEntry, sink: ScrcpyStreamSink): void { - entry.subscribers.delete(sink) - entry.waiting.delete(sink) - if (entry.subscribers.size > 0 || entry.closed || entry.idleTimer !== undefined) return - entry.idleTimer = setTimeout(() => { - entry.idleTimer = undefined - if (entry.subscribers.size === 0) void this.closeEntry(entry) - }, this.idleGraceMs) +import { ScrcpyVideoStreams as RuntimeStreams, buildScrcpyVideoServerArgs as serverArgs, + type ScrcpyVideoStreamsOptions as RuntimeOptions } from '../../packages/device-runtime/src/scrcpy-stream.ts' +import { SCRCPY_VERSION, resolveScrcpyAsset } from './scrcpy.ts' +export { ScrcpyVideoPacketParser } from '../../packages/device-runtime/src/scrcpy-stream.ts' +export type { VideoDevice, ScrcpyStreamSink, ScrcpyStreamStatus, ScrcpyVideoEvent } from '../../packages/device-runtime/src/scrcpy-stream.ts' +export type ScrcpyVideoStreamsOptions = Omit +const remoteServer = '/data/local/tmp/opengui-workbuddy-scrcpy-server.jar' +export function buildScrcpyVideoServerArgs(scid: string, serverPath = remoteServer): string[] { + return serverArgs(scid, serverPath, SCRCPY_VERSION) +} +export class ScrcpyVideoStreams extends RuntimeStreams { + constructor(options: ScrcpyVideoStreamsOptions) { + super({ ...options, asset: options.asset ?? resolveScrcpyAsset(), version: SCRCPY_VERSION, remoteServer }) } - - private async start(entry: StreamEntry, asset: ScrcpyAsset): Promise { - const signal = AbortSignal.any([entry.controller.signal, this.lifetime.signal]) - this.phase = 'downloading' - this.message = undefined - const installed = await this.installer.ensure(asset, signal, progress => { - this.phase = progress.phase - this.downloadedBytes = progress.downloadedBytes - this.broadcastText(entry, { type: 'install', phase: progress.phase, downloadedBytes: progress.downloadedBytes, totalBytes: progress.totalBytes }) - }) - this.phase = 'ready' - this.downloadedBytes = undefined - signal.throwIfAborted() - - const port = await this.freePort() - entry.port = port - const scid = (randomBytes(4).readUInt32BE(0) & 0x7fffffff).toString(16).padStart(8, '0') - const forward: OwnedForward = { serial: entry.device.serial, port, scid, kind: 'video-stream' } - await this.options.runAdb(['-s', entry.device.serial, 'push', installed.server, SCRCPY_REMOTE_SERVER], signal) - try { - await this.forwardRegistry.track(forward) - entry.forward = forward - await this.options.runAdb([ - '-s', entry.device.serial, 'forward', '--no-rebind', `tcp:${port}`, `localabstract:scrcpy_${scid}`, - ], signal) - } catch (error) { - await this.forwardRegistry.release(forward, this.options.runAdb).catch(() => false) - throw error - } - - const child = this.spawnImpl(this.options.adbPath(), [ - '-s', entry.device.serial, 'shell', ...buildScrcpyVideoServerArgs(scid), - ], { shell: false, windowsHide: true, stdio: ['ignore', 'ignore', 'pipe'] }) - entry.process = child - let stderr = '' - child.stderr?.on('data', (chunk: Buffer | string) => { stderr = `${stderr}${String(chunk)}`.slice(-2_000) }) - await waitForSpawn(child, signal) - const socket = await connectVideo(this.connectImpl, port, signal) - entry.socket = socket - const parser = new ScrcpyVideoPacketParser() - socket.on('data', chunk => { - try { - for (const event of parser.push(Buffer.from(chunk))) this.broadcastEvent(entry, event) - } catch (error) { - this.broadcastText(entry, { type: 'error', message: this.publicError(error) }) - void this.closeEntry(entry) - } - }) - const settled = new Promise((resolve, reject) => { - let socketCloseTimer: ReturnType | undefined - const rejectSocketClose = (): void => { - socketCloseTimer = setTimeout(() => { - reject(new Error(`scrcpy video socket closed unexpectedly${stderr.trim() ? `: ${stderr.trim()}` : ''}`)) - }, 150) - } - socket.once('error', reject) - socket.once('close', () => signal.aborted ? resolve() : rejectSocketClose()) - child.once('exit', (code, exitSignal) => { - if (socketCloseTimer !== undefined) clearTimeout(socketCloseTimer) - if (entry.controller.signal.aborted) resolve() - else reject(new Error(`scrcpy video server exited (code=${String(code)}, signal=${String(exitSignal)})${stderr.trim() ? `: ${stderr.trim()}` : ''}`)) - }) - signal.addEventListener('abort', () => { - if (socketCloseTimer !== undefined) clearTimeout(socketCloseTimer) - resolve() - }, { once: true }) - }) - await settled - } - - private broadcastEvent(entry: StreamEntry, event: ScrcpyVideoEvent): void { - if (event.type === 'packet') { - const frame = this.packetFrame(event) - if (event.config) { - entry.replay = [frame] - entry.replayBytes = frame.byteLength - } else if (event.key) { - entry.replay = [...entry.replay.filter(packet => (packet[0]! & 1) !== 0), frame] - entry.replayBytes = entry.replay.reduce((total, packet) => total + packet.byteLength, 0) - } else if (entry.replay.some(packet => (packet[0]! & 2) !== 0)) { - if (entry.replayBytes + frame.byteLength <= 8 * 1024 * 1024) { - entry.replay.push(frame) - entry.replayBytes += frame.byteLength - } else { - entry.replay = entry.replay.filter(packet => (packet[0]! & 1) !== 0) - entry.replayBytes = entry.replay.reduce((total, packet) => total + packet.byteLength, 0) - } - } - for (const sink of entry.subscribers) { - if (sink.bufferedBytes() > 1_000_000) { - entry.waiting.add(sink) - if (sink.bufferedBytes() > 2_000_000) sink.close(1013, 'slow_client') - continue - } - if (entry.waiting.has(sink)) { - if (!event.key) continue - sink.sendText(JSON.stringify({ type: 'reset' })) - for (const config of entry.replay.filter(packet => (packet[0]! & 1) !== 0)) sink.sendBinary(config) - entry.waiting.delete(sink) - } - sink.sendBinary(frame) - } - return - } - const text = JSON.stringify(event) - if (event.type === 'codec') entry.lastCodec = text - else { - entry.lastSession = text - entry.replay = [] - entry.replayBytes = 0 - } - for (const sink of entry.subscribers) sink.sendText(text) - } - - private packetFrame(event: Extract): Buffer { - const frame = Buffer.allocUnsafe(9 + event.data.length) - frame[0] = (event.config ? 1 : 0) | (event.key ? 2 : 0) - frame.writeBigUInt64BE(event.pts, 1) - event.data.copy(frame, 9) - return frame - } - - private broadcastText(entry: StreamEntry, value: unknown): void { - const text = JSON.stringify(value) - for (const sink of entry.subscribers) sink.sendText(text) - } - - private closeEntry(entry: StreamEntry): Promise { - if (entry.closing !== undefined) return entry.closing - if (entry.closed) return Promise.resolve() - entry.closed = true - const closing = (async (): Promise => { - if (entry.idleTimer !== undefined) clearTimeout(entry.idleTimer) - entry.controller.abort(new Error('opengui-workbuddy: embedded stream stopped')) - entry.socket?.destroy() - const child = entry.process - if (child !== undefined && child.exitCode === null) await terminateChild(child) - if (entry.forward !== undefined) await this.forwardRegistry.release(entry.forward, this.options.runAdb, 1_000) - if (this.entries.get(entry.device.id) === entry) this.entries.delete(entry.device.id) - for (const sink of entry.subscribers) sink.close(entry.closeCode, entry.closeReason) - entry.subscribers.clear() - })() - entry.closing = closing - return closing - } - - private publicError(_error: unknown): string { - return 'video_failed: encoder or device connection failed; reconnect the selected phone and retry video' - } -} - -async function terminateChild(child: ChildProcess): Promise { - if (child.exitCode !== null) return - const exited = new Promise(resolve => child.once('exit', () => resolve())) - if (!child.killed) child.kill('SIGTERM') - const graceful = await Promise.race([ - exited.then(() => true), - new Promise(resolve => setTimeout(() => resolve(false), 500)), - ]) - if (!graceful && child.exitCode === null) { - child.kill('SIGKILL') - await Promise.race([exited, new Promise(resolve => setTimeout(resolve, 250))]) - } -} - -async function availableTcpPort(): Promise { - return new Promise((resolvePort, rejectPort) => { - const server = createServer() - server.once('error', rejectPort) - server.listen(0, '127.0.0.1', () => { - const address = server.address() - const port = typeof address === 'object' && address !== null ? address.port : 0 - server.close(error => error === undefined ? resolvePort(port) : rejectPort(error)) - }) - }) -} - -async function waitForSpawn(child: ChildProcess, signal: AbortSignal): Promise { - await new Promise((resolve, reject) => { - const done = (error?: Error): void => { - child.off('spawn', onSpawn) - child.off('error', onError) - signal.removeEventListener('abort', onAbort) - error === undefined ? resolve() : reject(error) - } - const onSpawn = (): void => done() - const onError = (error: Error): void => done(error) - const onAbort = (): void => done(signal.reason instanceof Error ? signal.reason : new Error(String(signal.reason))) - child.once('spawn', onSpawn) - child.once('error', onError) - signal.addEventListener('abort', onAbort, { once: true }) - }) -} - -async function connectVideo(connectImpl: typeof connect, port: number, signal: AbortSignal): Promise { - const deadline = Date.now() + 10_000 - while (true) { - signal.throwIfAborted() - try { - const socket = await new Promise((resolve, reject) => { - const socket = connectImpl({ host: '127.0.0.1', port }) - const cleanup = (): void => { - socket.off('connect', onConnect) - socket.off('error', onError) - signal.removeEventListener('abort', onAbort) - } - const onConnect = (): void => { cleanup(); resolve(socket) } - const onError = (error: Error): void => { cleanup(); socket.destroy(); reject(error) } - const onAbort = (): void => { cleanup(); socket.destroy(); reject(signal.reason) } - socket.once('connect', onConnect) - socket.once('error', onError) - signal.addEventListener('abort', onAbort, { once: true }) - }) - try { - await waitForVideoData(socket, signal, Math.max(1, deadline - Date.now())) - return socket - } catch (error) { - socket.destroy() - throw error - } - } catch (error) { - if (signal.aborted || Date.now() >= deadline) throw error - await new Promise(resolve => setTimeout(resolve, 120)) - } - } -} - -/** ADB forward accepts TCP before the device abstract socket exists, then closes it. - * Treat a connection as ready only after scrcpy has produced its first bytes. */ -async function waitForVideoData(socket: Socket, signal: AbortSignal, timeoutMs: number): Promise { - await new Promise((resolve, reject) => { - const timeout = setTimeout(() => done(new Error('opengui-workbuddy: timed out waiting for scrcpy video data')), timeoutMs) - const done = (error?: Error): void => { - clearTimeout(timeout) - socket.off('readable', onReadable) - socket.off('end', onClose) - socket.off('close', onClose) - socket.off('error', onError) - signal.removeEventListener('abort', onAbort) - error === undefined ? resolve() : reject(error) - } - const onReadable = (): void => { - if (socket.readableLength > 0) done() - } - const onClose = (): void => done(new Error('opengui-workbuddy: scrcpy video socket not ready')) - const onError = (error: Error): void => done(error) - const onAbort = (): void => done(signal.reason instanceof Error ? signal.reason : new Error(String(signal.reason))) - socket.once('readable', onReadable) - socket.once('end', onClose) - socket.once('close', onClose) - socket.once('error', onError) - signal.addEventListener('abort', onAbort, { once: true }) - }) } diff --git a/workbuddy-plugin/src/service.ts b/workbuddy-plugin/src/service.ts index 2bc1cd8..304256b 100644 --- a/workbuddy-plugin/src/service.ts +++ b/workbuddy-plugin/src/service.ts @@ -1,3 +1,4 @@ +import { SessionRuntime } from '../../packages/device-runtime/src/session-runtime.ts' import { ViewerServer, type ViewerStreams } from './viewer.ts' import { ScrcpyVideoStreams } from './scrcpy-stream.ts' import { randomUUID } from 'node:crypto' @@ -416,8 +417,9 @@ export class WorkBuddyOpenGuiService { private readonly now: () => number private readonly host: WorkBuddyPhoneHost private readonly createSessionId: () => string - private readonly sessions = new Map() - private readonly locks = new Map() + private readonly runtime = new SessionRuntime() + private readonly sessions = this.runtime.sessions + private readonly locks = this.runtime.locks readonly viewers: ViewerServer private readonly wall: DeviceWallServer private disposed = false @@ -566,9 +568,8 @@ export class WorkBuddyOpenGuiService { task.actors.set(item.device.serial, item.actor) } this.host.assignTarget(item.actor, item.device.serial) - if (purpose === 'control') this.locks.set(item.device.serial, id) } - this.sessions.set(id, record) + this.runtime.register(record, purpose === 'control') this.renewLease(record) try { await this.wall.start() @@ -660,7 +661,7 @@ export class WorkBuddyOpenGuiService { delete action.confirmationRequestId delete action.hostContext try { - return await this.runPhoneOperation(sessionId, deviceId, signal, (item, combined) => this.host.act(item.actor, action, combined)) + return await this.runPhoneOperation(sessionId, deviceId, signal, (item, combined) => this.host.act(item.actor, action, combined), true) } catch (error) { item.resultUnknown ||= errorInfo(error).executionState === 'outcome_unknown' record.resultUnknown = record.devices.some(device => device.resultUnknown) @@ -773,12 +774,14 @@ export class WorkBuddyOpenGuiService { deviceId: string | undefined, signal: AbortSignal, operation: (item: SessionDevice, combined: AbortSignal) => Promise, + mutating = false, ): Promise { const record = this.requireActiveSession(sessionId) if (record.purpose === 'mirror') throw new Error('opengui: mirror-only sessions cannot capture model images or control phones') const item = this.resolveDevice(record, deviceId) const combined = AbortSignal.any([record.controller.signal, item.connectionController.signal, signal]) const connectionEpoch = item.connectionEpoch ?? 0 + let completed = false try { combined.throwIfAborted() this.viewers.assertReady(record.viewerId!) @@ -787,6 +790,7 @@ export class WorkBuddyOpenGuiService { record.task.operations.set(item.device.serial, operations + 1) this.renewLease(record) const value = await this.track(record, () => operation(item, combined)) + completed = true combined.throwIfAborted() if ((item.connectionEpoch ?? 0) !== connectionEpoch) { this.host.invalidate?.(item.actor) @@ -803,6 +807,7 @@ export class WorkBuddyOpenGuiService { item.needsObservation = true this.host.invalidate?.(item.actor) record.lastError = error instanceof Error ? error.message : String(error) + if (mutating && completed) throw new OpenGuiError('result_delivery_failed', record.lastError, 'outcome_unknown', 'observe') throw error } } @@ -829,51 +834,24 @@ export class WorkBuddyOpenGuiService { } } - private requireSession(sessionId: string): SessionRecord { - const record = this.sessions.get(sessionId) - if (record === undefined) throw new Error('opengui: unknown sessionId') - return record - } + private requireSession(sessionId: string): SessionRecord { return this.runtime.requireSession(sessionId) } - private requireActiveSession(sessionId: string): SessionRecord { - const record = this.requireSession(sessionId) - if (record.state !== 'active') throw new Error(`opengui: session is ${record.state}`) - return record - } + private requireActiveSession(sessionId: string): SessionRecord { return this.runtime.requireActiveSession(sessionId) } private resolveDevice(record: SessionRecord, deviceId: string | undefined): SessionDevice { - if (deviceId === undefined) { - if (record.devices.length !== 1) throw new Error('opengui: deviceId is required for a multi-device session') - return record.devices[0]! - } - const item = record.devices.find(candidate => candidate.device.id === deviceId) - if (item === undefined) throw new Error('opengui: deviceId is not locked by this session') - return item + return this.runtime.resolveDevice(record, deviceId) } - private release(record: SessionRecord): void { - for (const item of record.devices) { - if (this.locks.get(item.device.serial) === record.id) this.locks.delete(item.device.serial) - } - } + private release(record: SessionRecord): void { this.runtime.release(record) } - private async track(record: SessionRecord, operation: () => Promise): Promise { - const pending = Promise.resolve().then(() => { - record.controller.signal.throwIfAborted() - return operation() - }) - record.pending.add(pending) - try { return await pending } finally { record.pending.delete(pending) } + private track(record: SessionRecord, operation: () => Promise): Promise { + return this.runtime.track(record, operation) } /** Keep leases until in-flight work and owned resource cleanup have both drained. */ private cleanup(record: SessionRecord): Promise { clearTimeout(record.leaseTimer) - record.cleanup ??= (async () => { - await Promise.allSettled([...record.pending]) - await this.releaseDeviceResources(record) - this.release(record) - })() + record.cleanup ??= this.runtime.drain(record, () => this.releaseDeviceResources(record)) return record.cleanup } diff --git a/workbuddy-plugin/tests/runtime-contract.spec.ts b/workbuddy-plugin/tests/runtime-contract.spec.ts index 04493b7..97d139c 100644 --- a/workbuddy-plugin/tests/runtime-contract.spec.ts +++ b/workbuddy-plugin/tests/runtime-contract.spec.ts @@ -2,3 +2,5 @@ import { it } from 'vitest' import { PhoneController } from '../src/phone-controller.ts' import { controllerContract } from '../../packages/device-runtime/tests/controller-contract.ts' controllerContract(it, options => new PhoneController(options)) +import { sessionContract } from '../../packages/device-runtime/tests/session-contract.ts' +sessionContract(it)