diff --git a/apps/mobile/src/features/devices/device-stream.browser.test.ts b/apps/mobile/src/features/devices/device-stream.browser.test.ts index 5aeefdac0203..24766242b37e 100644 --- a/apps/mobile/src/features/devices/device-stream.browser.test.ts +++ b/apps/mobile/src/features/devices/device-stream.browser.test.ts @@ -38,6 +38,7 @@ async function setup() { static OPEN = 1; readyState = 1; onopen: (() => void) | null = null; + onmessage: ((event: { data: ArrayBuffer }) => void) | null = null; close = vi.fn(); send = vi.fn(); constructor() { @@ -86,6 +87,8 @@ it("bridges shared first-frame readiness, image failure, and a successful fresh const { elements, messages, sockets, configuration } = await setup(); sockets[0]!.onopen?.(); expect(messages()).not.toContainEqual({ type: "status", status: "streaming" }); + expect(messages()).not.toContainEqual({ type: "input", connected: true }); + sockets[0]!.onmessage?.({ data: new Uint8Array([0x83]).buffer }); expect(messages()).toContainEqual({ type: "input", connected: true }); const image = elements.find((element) => element.tag === "img")!; image.naturalWidth = 400; diff --git a/apps/server/src/device/DeviceToolchain.ts b/apps/server/src/device/DeviceToolchain.ts index f30960ed2212..af590787e4ee 100644 --- a/apps/server/src/device/DeviceToolchain.ts +++ b/apps/server/src/device/DeviceToolchain.ts @@ -26,9 +26,9 @@ import * as Semaphore from "effect/Semaphore"; import * as ProcessRunner from "../processRunner.ts"; const DEVICE_HUB_PACKAGE = "expo-device-hub"; -export const DEVICE_HUB_VERSION = "0.12.0"; +export const DEVICE_HUB_VERSION = "0.15.3"; const AGENT_DEVICE_PACKAGE = "agent-device"; -export const AGENT_DEVICE_VERSION = "0.21.12"; +export const AGENT_DEVICE_VERSION = "0.21.23"; const INSTALL_TIMEOUT = Duration.minutes(10); const installLock = Semaphore.makeUnsafe(1); diff --git a/apps/web/src/components/device/DeviceStreamView.test.tsx b/apps/web/src/components/device/DeviceStreamView.test.tsx index 7e0392340129..a7f4079b5683 100644 --- a/apps/web/src/components/device/DeviceStreamView.test.tsx +++ b/apps/web/src/components/device/DeviceStreamView.test.tsx @@ -1,4 +1,4 @@ -import { act, useSyncExternalStore } from "react"; +import { act, useEffect, useSyncExternalStore } from "react"; import { create, type ReactTestRenderer } from "react-test-renderer"; import { EnvironmentId } from "@t3tools/contracts"; import { afterEach, beforeEach, expect, it, vi } from "vite-plus/test"; @@ -20,6 +20,20 @@ vi.mock("~/state/device", () => ({ useDeviceHubAccess: () => useSyncExternalStore(accessStore.subscribe, () => accessStore.value), refreshDeviceHubAccess: () => accessStore.refresh(), })); +// Replace GPU allocation while keeping the real React and stream lifecycles. +let viewerMounts = 0; +let viewerUnmounts = 0; +vi.mock("./DevicePhoneViewport", () => ({ + DevicePhoneViewport: function Viewer() { + useEffect(() => { + viewerMounts++; + return () => { + viewerUnmounts++; + }; + }, []); + return null; + }, +})); import { DeviceStreamView } from "./DeviceStreamView"; class Image extends EventTarget { @@ -34,6 +48,8 @@ let renderer: ReactTestRenderer | undefined; let primes = 0; beforeEach(() => { primes = 0; + viewerMounts = 0; + viewerUnmounts = 0; }); afterEach(async () => { await act(async () => renderer?.unmount()); @@ -42,18 +58,57 @@ afterEach(async () => { vi.unstubAllGlobals(); }); -async function setup() { +async function setup(h264 = false) { vi.useFakeTimers(); + vi.stubGlobal("window", globalThis); vi.stubGlobal("IS_REACT_ACT_ENVIRONMENT", true); - vi.stubGlobal("fetch", () => { + let videoBody: ReadableStreamDefaultController | undefined; + let output: VideoFrameOutputCallback | undefined; + if (h264) { + vi.stubGlobal( + "VideoDecoder", + class { + static isConfigSupported = async () => ({ supported: true }); + state = "unconfigured"; + constructor(callbacks: VideoDecoderInit) { + output = callbacks.output; + } + configure() { + this.state = "configured"; + } + close() { + this.state = "closed"; + } + }, + ); + vi.stubGlobal("EncodedVideoChunk", vi.fn()); + } + vi.stubGlobal("fetch", (url: string, init: RequestInit) => { + if (url.endsWith("stream.avcc")) { + return Promise.resolve( + new Response( + new ReadableStream({ + start(controller) { + videoBody = controller; + init.signal?.addEventListener("abort", () => controller.error(new Error("aborted"))); + }, + }), + ), + ); + } primes++; return Promise.resolve(new Response("prime")); }); + const sockets: Array<{ onmessage?: (event: { data: ArrayBuffer }) => void }> = []; vi.stubGlobal( "WebSocket", class { static OPEN = 1; readyState = 1; + constructor() { + sockets.push(this); + } + onmessage?: (event: { data: ArrayBuffer }) => void; send() {} close() {} }, @@ -73,6 +128,7 @@ async function setup() { deviceId="test" platform="ios" visible={visible} + allowPhoneView={h264} /> ); await act(async () => { @@ -84,13 +140,30 @@ async function setup() { return image; } return { + getContext: () => ({ drawImage() {} }), style: { setProperty() {} }, getBoundingClientRect: () => ({ width: 400, height: 800 }), }; }, }); }); - return { images, view }; + return { + images, + view, + configure() { + const json = new TextEncoder().encode( + JSON.stringify({ width: 400, height: 800, orientation: "portrait" }), + ); + const packet = new Uint8Array(1 + json.length); + packet[0] = 0x82; + packet.set(json, 1); + sockets[0]?.onmessage?.({ data: packet.buffer }); + videoBody?.enqueue(new Uint8Array([0, 0, 0, 5, 1, 1, 0x64, 0, 0x1f])); + }, + frame() { + output?.({ displayWidth: 400, displayHeight: 800, close() {} } as VideoFrame); + }, + }; } it("removes MJPEG requests while hidden and reconnects when shown", async () => { @@ -135,3 +208,25 @@ it("starts exactly one new stream per Reconnect press", async () => { await act(async () => renderer!.root.findByType("button").props.onClick()); expect(primes).toBe(2); }); + +it("keeps the 3D viewer mounted across iOS video recovery and releases it when hidden", async () => { + const { configure, frame, view } = await setup(true); + await act(async () => configure()); + expect(viewerMounts).toBe(0); + await act(async () => frame()); + expect(viewerMounts).toBe(1); + await act(async () => { + await vi.advanceTimersByTimeAsync(15_000); + }); + expect(viewerUnmounts).toBe(0); + expect(viewerMounts).toBe(1); + await act(async () => { + await vi.advanceTimersByTimeAsync(1_000); + configure(); + }); + await act(async () => frame()); + expect(viewerMounts).toBe(1); + await act(async () => renderer!.update(view(false))); + expect(viewerUnmounts).toBe(1); + expect(vi.getTimerCount()).toBe(0); +}); diff --git a/apps/web/src/components/device/DeviceStreamView.tsx b/apps/web/src/components/device/DeviceStreamView.tsx index 25c4c39dae2f..8db311f2fe14 100644 --- a/apps/web/src/components/device/DeviceStreamView.tsx +++ b/apps/web/src/components/device/DeviceStreamView.tsx @@ -94,6 +94,7 @@ export function DeviceStreamView(props: { const canvasRef = useRef(null); const clientRef = useRef(null); const [status, setStatus] = useState("connecting"); + const [hasFrame, setHasFrame] = useState(false); const [detail, setDetail] = useState(undefined); const [showRestartNotice, setShowRestartNotice] = useState(false); const [screen, setScreen] = useState(null); @@ -123,6 +124,7 @@ export function DeviceStreamView(props: { onDuoUnavailable: onPhoneUnavailable, onStatus: (next, nextDetail) => { setStatus(next); + if (next === "streaming") setHasFrame(true); setDetail(nextDetail); if (next !== "connecting") setShowRestartNotice(false); }, @@ -150,6 +152,8 @@ export function DeviceStreamView(props: { }, ); clientRef.current = client; + setPhoneUnavailable(false); + setHasFrame(false); setMjpegUrl(null); setInputState({ connected: false }); client.start(); @@ -188,16 +192,13 @@ export function DeviceStreamView(props: { return w / h; }, [props.platform, screen]); - // Android restarts its encoder when a fold changes the framebuffer size. - // Keep the last decoded frame and viewer mounted while the next keyframe arrives. - const retainingAndroidFrame = - props.platform === "android" && - status === "connecting" && - inputState.connected && - screen !== null; + // Encoder restarts and video reconnects retain the decoded frame and viewer + // while input is still connected and the next keyframe is on its way. + const retainingFrame = + status === "connecting" && hasFrame && inputState.connected && screen !== null; const showPhone = props.allowPhoneView && - (status === "streaming" || retainingAndroidFrame) && + (status === "streaming" || retainingFrame) && props.visible && presentation === "phone" && !phoneUnavailable && @@ -205,10 +206,10 @@ export function DeviceStreamView(props: { !props.axOverlay && (!isDuo || screen?.supportsHingeAngle === true); useEffect(() => { - if (!retainingAndroidFrame || !showPhone) return; + if (!retainingFrame || !showPhone) return; const timeout = window.setTimeout(() => setShowRestartNotice(true), 2_000); return () => window.clearTimeout(timeout); - }, [retainingAndroidFrame, showPhone]); + }, [retainingFrame, showPhone]); const controlsInset = props.renderControls && !showPhone ? CONTROLS_RAIL_WIDTH : 0; // The frame is the largest box at `aspect` that fits the container, so a @@ -532,14 +533,14 @@ export function DeviceStreamView(props: { ) : null} - {retainingAndroidFrame && showPhone && showRestartNotice ? ( + {retainingFrame && showPhone && showRestartNotice ? (
Waiting for device video…
) : null} - {status !== "streaming" && !(retainingAndroidFrame && showPhone) ? ( + {status !== "streaming" && !(retainingFrame && showPhone) ? (
diff --git a/packages/client-runtime/src/device/stream.test.ts b/packages/client-runtime/src/device/stream.test.ts index 7ea5aca27abf..f29484eb00b4 100644 --- a/packages/client-runtime/src/device/stream.test.ts +++ b/packages/client-runtime/src/device/stream.test.ts @@ -385,10 +385,14 @@ function recoveryFixture(platform: "ios" | "android" = "ios", preferMjpeg = true configure() { this.state = "configured"; } + reset = vi.fn(() => { + this.state = "unconfigured"; + this.decodeQueueSize = 0; + }); close() { this.state = "closed"; } - decode() {} + decode = vi.fn(); readonly callbacks: { output: (frame: VideoFrame) => void; error: () => void }; constructor(callbacks: { output: (frame: VideoFrame) => void; error: () => void }) { this.callbacks = callbacks; @@ -480,6 +484,93 @@ describe("shared device stream readiness and recovery", () => { vi.unstubAllGlobals(); }); + it("waits for iOS input admission and respects native input loss and recovery", async () => { + const { client, sockets, events } = recoveryFixture(); + client.start(); + await vi.advanceTimersByTimeAsync(0); + const socket = sockets[0]!; + socket.onopen?.(); + expect(events.onInputConnected).not.toHaveBeenCalledWith(true); + socket.onmessage?.({ data: new Uint8Array([0x83]).buffer }); + expect(events.onInputConnected).toHaveBeenLastCalledWith(true); + const config = (inputUnavailable: boolean) => { + const json = new TextEncoder().encode( + JSON.stringify({ width: 400, height: 800, orientation: "portrait", inputUnavailable }), + ); + const packet = new Uint8Array(1 + json.length); + packet[0] = 0x82; + packet.set(json, 1); + socket.onmessage?.({ data: packet.buffer }); + }; + config(true); + expect(events.onInputConnected).toHaveBeenLastCalledWith( + false, + "Simulator input is unavailable", + ); + socket.send.mockClear(); + client.sendTouch("begin", 0.2, 0.3); + client.pressButton("home"); + expect(socket.send).not.toHaveBeenCalled(); + config(false); + expect(events.onInputConnected).toHaveBeenLastCalledWith(true, undefined); + client.sendTouch("begin", 0.2, 0.3); + expect(socket.send).toHaveBeenCalledOnce(); + client.stop(); + expect(vi.getTimerCount()).toBe(0); + }); + + it("drops an iOS decode backlog and resumes at a keyframe without switching to MJPEG", async () => { + const { client, videoBodies, decoders, events, decodedFrame } = recoveryFixture("ios", false); + client.start(); + await vi.advanceTimersByTimeAsync(0); + const body = videoBodies[0]!; + body.enqueue(new Uint8Array([...envelope(1, [1, 0x64, 0, 0x1f]), ...envelope(2, [1])])); + await vi.advanceTimersByTimeAsync(0); + decodedFrame(); + const decoder = decoders[0]!; + expect(decoder.decode).toHaveBeenCalledOnce(); + decoder.decodeQueueSize = 9; + body.enqueue(new Uint8Array([...envelope(3, [2]), ...envelope(3, [3])])); + await vi.advanceTimersByTimeAsync(0); + expect(decoder.reset).toHaveBeenCalledOnce(); + expect(decoder.decode).toHaveBeenCalledOnce(); + body.enqueue(new Uint8Array([...envelope(2, [4]), ...envelope(3, [5])])); + await vi.advanceTimersByTimeAsync(0); + expect(decoder.decode).toHaveBeenCalledTimes(3); + expect(events.onMjpegFallback).not.toHaveBeenCalled(); + expect(events.onStatus).toHaveBeenLastCalledWith("streaming", undefined); + client.stop(); + }); + + it("does not repaint a delayed JPEG seed over newer decoded video", async () => { + const { client, videoBodies, drawImage, decodedFrame } = recoveryFixture("ios", false); + let resolveBitmap!: (bitmap: ImageBitmap) => void; + vi.stubGlobal( + "createImageBitmap", + () => + new Promise((resolve) => { + resolveBitmap = resolve; + }), + ); + client.start(); + await vi.advanceTimersByTimeAsync(0); + videoBodies[0]!.enqueue( + new Uint8Array([ + ...envelope(4, [1]), + ...envelope(1, [1, 0x64, 0, 0x1f]), + ...envelope(2, [2]), + ]), + ); + await vi.advanceTimersByTimeAsync(0); + decodedFrame(); + const close = vi.fn(); + resolveBitmap({ width: 400, height: 800, close } as unknown as ImageBitmap); + await vi.advanceTimersByTimeAsync(0); + expect(drawImage).toHaveBeenCalledOnce(); + expect(close).toHaveBeenCalledOnce(); + client.stop(); + }); + it("waits for actual MJPEG dimensions without requiring a multipart load event", async () => { const { client, events, image } = recoveryFixture(); client.start(); @@ -713,22 +804,35 @@ describe("shared device stream readiness and recovery", () => { }); }); -it("reports a stalled AVCC body after its initial image instead of leaving a frozen streaming state", async () => { - const { client, videoBodies, events, signals } = recoveryFixture("ios", false); +it("reopens a stalled AVCC body without disconnecting input, then receives fresh video", async () => { + const { client, videoBodies, events, signals, sockets, decodedFrame } = recoveryFixture( + "ios", + false, + ); const close = vi.fn(); vi.stubGlobal("createImageBitmap", () => Promise.resolve({ width: 400, height: 800, close })); try { client.start(); await vi.advanceTimersByTimeAsync(0); + sockets[0]!.onmessage?.({ data: new Uint8Array([0x83]).buffer }); videoBodies[0]!.enqueue(new Uint8Array(envelope(4, [1]))); await vi.advanceTimersByTimeAsync(0); expect(events.onStatus).toHaveBeenLastCalledWith("streaming", undefined); await vi.advanceTimersByTimeAsync(15_000); expect(events.onStatus).toHaveBeenLastCalledWith( - "error", + "connecting", expect.stringContaining("stopped receiving video"), ); - expect(signals.every((signal) => signal.aborted)).toBe(true); + expect(signals[1]!.aborted).toBe(true); + expect(sockets[0]!.close).not.toHaveBeenCalled(); + expect(events.onInputConnected).toHaveBeenLastCalledWith(true); + await vi.advanceTimersByTimeAsync(1_000); + expect(videoBodies).toHaveLength(2); + videoBodies[1]!.enqueue(new Uint8Array(envelope(1, [1, 0x64, 0, 0x1f]))); + await vi.advanceTimersByTimeAsync(0); + decodedFrame(); + expect(events.onStatus).toHaveBeenLastCalledWith("streaming", undefined); + client.stop(); expect(vi.getTimerCount()).toBe(0); } finally { client.stop(); diff --git a/packages/client-runtime/src/device/stream.ts b/packages/client-runtime/src/device/stream.ts index 6b1ff6158879..27ab7c93b416 100644 --- a/packages/client-runtime/src/device/stream.ts +++ b/packages/client-runtime/src/device/stream.ts @@ -55,6 +55,7 @@ const screenConfigSchema = Schema.Struct({ "landscape_left", "landscape_right", ]), + inputUnavailable: Schema.optionalKey(Schema.Boolean), screenId: Schema.optionalKey(Schema.Number), supportsHingeAngle: Schema.optionalKey(Schema.Boolean), supportsPhysicalOrientation: Schema.optionalKey(Schema.Boolean), @@ -130,6 +131,7 @@ const IOS_MSG_ORIENTATION = 0x07; const IOS_MSG_HARDWARE_KEYBOARD = 0x0d; // helper -> browser. const IOS_TAG_SCREEN_CONFIG = 0x82; +const IOS_TAG_INPUT_ADMITTED = 0x83; const encoder = new TextEncoder(); const decoder = new TextDecoder(); @@ -364,6 +366,9 @@ export function createDeviceStreamClient( const retryTimers = new Map<"video" | "input", ReturnType>(); let primeController: AbortController | null = null; let videoDecoder: VideoDecoder | null = null; + let decoderConfig: VideoDecoderConfig | null = null; + let decodedVideoGeneration: number | null = null; + let iosInputUnavailable = false; let timestamp = 0; let awaitingKeyframe = true; let screen: DeviceScreenSize | null = null; @@ -517,6 +522,7 @@ export function createDeviceStreamClient( // Already closed. } videoDecoder = null; + decoderConfig = null; awaitingKeyframe = true; }; @@ -537,8 +543,10 @@ export function createDeviceStreamClient( if ( videoDecoder === decoder && (platform !== "ios" || feedGeneration === videoGeneration) - ) + ) { + decodedVideoGeneration = feedGeneration; paint(frame, frame.displayWidth, frame.displayHeight); + } } finally { frame.close(); } @@ -572,6 +580,7 @@ export function createDeviceStreamClient( try { if (!videoDecoder || videoDecoder.state === "closed") videoDecoder = makeDecoder(); videoDecoder.configure(full); + decoderConfig = full; return true; } catch (cause) { if (platform === "android") fail(`Video decoder: ${(cause as Error).message}`); @@ -586,8 +595,22 @@ export function createDeviceStreamClient( awaitingKeyframe = false; } if (videoDecoder.decodeQueueSize > SOFT_DECODE_QUEUE) { - recoverDecoder(); - return; + if (platform !== "ios" || !decoderConfig) { + recoverDecoder(); + return; + } + // A slow viewer can accumulate healthy H.264 frames. Drop that backlog + // and resume at an IDR instead of permanently giving up its 3D view. + try { + videoDecoder.reset(); + videoDecoder.configure(decoderConfig); + awaitingKeyframe = true; + if (!isKey) return; + awaitingKeyframe = false; + } catch { + recoverDecoder(); + return; + } } try { videoDecoder.decode( @@ -659,7 +682,9 @@ export function createDeviceStreamClient( for (;;) { // An AVCC body can stay open after its helper stops producing frames. const timer = setTimeout(() => { - if (isCurrent()) fail("Device stream stopped receiving video. Reconnect to try again."); + if (!isCurrent()) return; + retryDetail = "Device stream stopped receiving video. Reconnecting…"; + videoController.abort(); }, FIRST_FRAME_TIMEOUT_MS); videoReadTimer = timer; let result: ReadableStreamReadResult; @@ -678,7 +703,8 @@ export function createDeviceStreamClient( void createImageBitmap(new Blob([chunk.payload as BlobPart], { type: "image/jpeg" })) .then((bitmap) => { try { - if (isCurrent()) paint(bitmap, bitmap.width, bitmap.height); + if (isCurrent() && decodedVideoGeneration !== feedGeneration) + paint(bitmap, bitmap.width, bitmap.height); } finally { bitmap.close(); } @@ -714,7 +740,7 @@ export function createDeviceStreamClient( } } catch (cause) { if (!isCurrent()) return; - retryDetail = (cause as Error).message; + retryDetail ??= (cause as Error).message; } if (isCurrent()) { controller = null; @@ -799,13 +825,17 @@ export function createDeviceStreamClient( ws.onopen = () => { if (stopped || socket !== ws) return; ws.send(taggedJson(IOS_MSG_HARDWARE_KEYBOARD, { enabled: false })); - events.onInputConnected(true); }; ws.onmessage = (event) => { if (stopped || socket !== ws) return; if (!(event.data instanceof ArrayBuffer)) return; const bytes = new Uint8Array(event.data); if (socket !== ws || stopped || bytes.length < 1) return; + if (bytes[0] === IOS_TAG_INPUT_ADMITTED) { + iosInputUnavailable = false; + events.onInputConnected(true); + return; + } try { const payload: unknown = JSON.parse(decoder.decode(bytes.subarray(1))); if (bytes[0] === 0x90) { @@ -814,6 +844,11 @@ export function createDeviceStreamClient( } else if (bytes[0] === IOS_TAG_SCREEN_CONFIG) { const config = decodeScreenConfig(payload); if (Option.isSome(config)) { + iosInputUnavailable = config.value.inputUnavailable === true; + events.onInputConnected( + !iosInputUnavailable, + iosInputUnavailable ? "Simulator input is unavailable" : undefined, + ); const previous = screen; screen = config.value; if (screen.hingePose && screen.hingePose !== previous?.hingePose) @@ -944,6 +979,7 @@ export function createDeviceStreamClient( stopped = false; generation++; configuring = false; + iosInputUnavailable = false; connecting(); if (platform === "ios") { if (!target.videoOnly) void connectIosInput(); @@ -986,7 +1022,8 @@ export function createDeviceStreamClient( }; const send = (payload: Uint8Array | string) => { - if (!stopped && socket?.readyState === WebSocket.OPEN) socket.send(payload); + if (!stopped && !iosInputUnavailable && socket?.readyState === WebSocket.OPEN) + socket.send(payload); }; const rawPoint = (x: number, y: number) => {