From 7074589c493fdf39786282c39f7986ec6338c3fa Mon Sep 17 00:00:00 2001 From: "claude[bot]" <41898282+claude[bot]@users.noreply.github.com> Date: Wed, 8 Jul 2026 15:44:53 +0000 Subject: [PATCH 1/4] fix: stream auto-reconnect and preserve scene entities during disconnect MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Add retryStream utility with exponential backoff (1s–30s) - Wrap streamEntityChanges and streamSceneChanges in useDrawService with retryStream so the draw backend reconnects automatically after errors - Guard entity cleanup in useGeometries and usePointclouds: when all queries disappear (transient disconnect), preserve existing ECS entities so the scene doesn't go blank; only destroy entities when the partID changes or other queries remain active Co-authored-by: Devin T. Currie --- .changeset/fix-scene-drops.md | 5 + src/lib/__tests__/retry-stream.spec.ts | 154 ++++++++++++++++++ src/lib/hooks/useGeometries.svelte.ts | 24 ++- src/lib/hooks/usePointclouds.svelte.ts | 16 +- .../DrawService/useDrawService.svelte.ts | 52 +++--- src/lib/retry-stream.ts | 51 ++++++ 6 files changed, 263 insertions(+), 39 deletions(-) create mode 100644 .changeset/fix-scene-drops.md create mode 100644 src/lib/__tests__/retry-stream.spec.ts create mode 100644 src/lib/retry-stream.ts diff --git a/.changeset/fix-scene-drops.md b/.changeset/fix-scene-drops.md new file mode 100644 index 000000000..8bcbbbdae --- /dev/null +++ b/.changeset/fix-scene-drops.md @@ -0,0 +1,5 @@ +--- +'@viamrobotics/motion-tools': patch +--- + +fix: stream auto-reconnect for draw service and preserve scene entities during disconnect diff --git a/src/lib/__tests__/retry-stream.spec.ts b/src/lib/__tests__/retry-stream.spec.ts new file mode 100644 index 000000000..dcaa7347c --- /dev/null +++ b/src/lib/__tests__/retry-stream.spec.ts @@ -0,0 +1,154 @@ +import { describe, expect, it, vi } from 'vitest' + +import { retryStream } from '../retry-stream' + +describe('retryStream', () => { + it('calls run and resolves when run succeeds', async () => { + const run = vi.fn().mockResolvedValue(undefined) + const controller = new AbortController() + + // run resolves once, retryStream will call it again — abort after first call + run.mockImplementation(async () => { + controller.abort() + }) + + await retryStream(run, controller.signal) + + expect(run).toHaveBeenCalledTimes(1) + }) + + it('retries when run throws', async () => { + vi.useFakeTimers() + + const controller = new AbortController() + let callCount = 0 + + const run = vi.fn().mockImplementation(async () => { + callCount++ + if (callCount < 3) { + throw new Error('stream error') + } + controller.abort() + }) + + const promise = retryStream(run, controller.signal) + // Advance through the backoff delays + await vi.advanceTimersByTimeAsync(1_000) + await vi.advanceTimersByTimeAsync(2_000) + + await promise + + expect(run).toHaveBeenCalledTimes(3) + + vi.useRealTimers() + }) + + it('stops retrying when signal is aborted', async () => { + vi.useFakeTimers() + + const controller = new AbortController() + const run = vi.fn().mockRejectedValue(new Error('stream error')) + const onRetry = vi.fn() + + const promise = retryStream(run, controller.signal, onRetry) + + // First call fails immediately, then waits for backoff + await vi.advanceTimersByTimeAsync(0) + expect(run).toHaveBeenCalledTimes(1) + + // Abort during backoff wait + controller.abort() + await vi.advanceTimersByTimeAsync(1_000) + + await promise + + // Should have called onRetry once, but not retried run + expect(onRetry).toHaveBeenCalledTimes(1) + expect(run).toHaveBeenCalledTimes(1) + + vi.useRealTimers() + }) + + it('calls onRetry with the current delay', async () => { + vi.useFakeTimers() + + const controller = new AbortController() + let callCount = 0 + + const run = vi.fn().mockImplementation(async () => { + callCount++ + if (callCount < 3) { + throw new Error('stream error') + } + controller.abort() + }) + + const onRetry = vi.fn() + const promise = retryStream(run, controller.signal, onRetry) + + await vi.advanceTimersByTimeAsync(1_000) + await vi.advanceTimersByTimeAsync(2_000) + + await promise + + expect(onRetry).toHaveBeenCalledTimes(2) + expect(onRetry).toHaveBeenNthCalledWith(1, 1_000) + expect(onRetry).toHaveBeenNthCalledWith(2, 2_000) + + vi.useRealTimers() + }) + + it('does not call onRetry and restarts immediately on clean stream end', async () => { + const controller = new AbortController() + let callCount = 0 + + const run = vi.fn().mockImplementation(async () => { + callCount++ + if (callCount === 1) return // clean end — server closed the stream + controller.abort() + }) + + const onRetry = vi.fn() + await retryStream(run, controller.signal, onRetry) + + expect(run).toHaveBeenCalledTimes(2) + expect(onRetry).not.toHaveBeenCalled() + }) + + it('resets delay after a successful run', async () => { + vi.useFakeTimers() + + const controller = new AbortController() + let callCount = 0 + + const run = vi.fn().mockImplementation(async () => { + callCount++ + // First call: fail + if (callCount === 1) throw new Error('fail') + // Second call: succeed (stream ended cleanly) + if (callCount === 2) return + // Third call: fail + if (callCount === 3) throw new Error('fail') + // Fourth call: abort + controller.abort() + }) + + const onRetry = vi.fn() + const promise = retryStream(run, controller.signal, onRetry) + + // First failure + 1s backoff + await vi.advanceTimersByTimeAsync(1_000) + // Second call succeeds, delay resets. Third call fails, should use 1s again + await vi.advanceTimersByTimeAsync(1_000) + // Fourth call - abort + await vi.advanceTimersByTimeAsync(2_000) + + await promise + + // Both retries should have used 1000ms (reset after success) + expect(onRetry).toHaveBeenNthCalledWith(1, 1_000) + expect(onRetry).toHaveBeenNthCalledWith(2, 1_000) + + vi.useRealTimers() + }) +}) diff --git a/src/lib/hooks/useGeometries.svelte.ts b/src/lib/hooks/useGeometries.svelte.ts index 2644d485c..37dad9fd0 100644 --- a/src/lib/hooks/useGeometries.svelte.ts +++ b/src/lib/hooks/useGeometries.svelte.ts @@ -222,19 +222,27 @@ export const provideGeometries = (partID: () => string) => { }) } - // Clean up owners whose queries disappeared entirely + // Clean up owners whose queries disappeared entirely. + // Guard: if ALL queries are gone (activeQueryKeys empty), the machine is likely + // temporarily disconnected — preserve entities so they reappear on reconnect. + // Only destroy when the partID changed (old-partID entities) or other queries + // are still active (connected machine, resource legitimately removed). + const anyQueriesActive = activeQueryKeys.size > 0 for (const [queryKey, keys] of queryEntityKeys) { if (!activeQueryKeys.has(queryKey)) { - for (const key of keys) { - const entity = entities.get(key) - if (entity && world.has(entity)) { - entity.destroy() + const queryPartID = queryKey.split(':')[0]! + if (queryPartID !== currentPartID || anyQueriesActive) { + for (const key of keys) { + const entity = entities.get(key) + if (entity && world.has(entity)) { + entity.destroy() + } + + entities.delete(key) } - entities.delete(key) + queryEntityKeys.delete(queryKey) } - - queryEntityKeys.delete(queryKey) } } }) diff --git a/src/lib/hooks/usePointclouds.svelte.ts b/src/lib/hooks/usePointclouds.svelte.ts index bfc851e5f..6a7fb15b5 100644 --- a/src/lib/hooks/usePointclouds.svelte.ts +++ b/src/lib/hooks/usePointclouds.svelte.ts @@ -188,13 +188,21 @@ export const providePointclouds = (partID: () => string) => { }) } - // clean up queries that disappeared entirely + // clean up queries that disappeared entirely. + // Guard: if ALL queries are gone (activeQueryKeys empty), the machine is likely + // temporarily disconnected — preserve entities so they reappear on reconnect. + // Only destroy when the partID changed (old-partID entities) or other queries + // are still active (connected machine, camera legitimately removed). + const anyQueriesActive = activeQueryKeys.size > 0 for (const [queryKey, entity] of entities) { if (!activeQueryKeys.has(queryKey)) { - if (world.has(entity)) { - entity.destroy() + const queryPartID = queryKey.split(':')[0]! + if (queryPartID !== currentPartID || anyQueriesActive) { + if (world.has(entity)) { + entity.destroy() + } + entities.delete(queryKey) } - entities.delete(queryKey) } } }) diff --git a/src/lib/plugins/DrawService/useDrawService.svelte.ts b/src/lib/plugins/DrawService/useDrawService.svelte.ts index d315fb4b4..9aa6b2e7c 100644 --- a/src/lib/plugins/DrawService/useDrawService.svelte.ts +++ b/src/lib/plugins/DrawService/useDrawService.svelte.ts @@ -22,6 +22,7 @@ import { } from '$lib/draw' import { hierarchy, traits, useWorld } from '$lib/ecs' import { useCameraControls } from '$lib/hooks/useControls.svelte' +import { retryStream } from '$lib/retry-stream' import { createServerRelationships } from './serverRelationships' import { useDrawConnectionConfig } from './useDrawConnectionConfig.svelte' @@ -304,33 +305,34 @@ export function provideDrawService() { } const streamEntityChanges = async (client: Client, signal: AbortSignal) => { - try { - for await (const response of client.streamEntityChanges({}, { signal })) { - connectionStatus = ConnectionStatus.CONNECTED - - const { entity } = response - if (!entity.case) continue - - const uuid = UuidTool.toString([...(entity.value.uuid ?? [])]) - pendingEvents.push({ - uuid, - changeType: response.changeType, - entity, - updatedFields: response.updatedFields, - }) - scheduleFlush() - } - } catch (error) { - if (!signal.aborted) { - console.error('Draw service entity stream error:', error) + await retryStream( + async (sig) => { + for await (const response of client.streamEntityChanges({}, { signal: sig })) { + connectionStatus = ConnectionStatus.CONNECTED + + const { entity } = response + if (!entity.case) continue + + const uuid = UuidTool.toString([...(entity.value.uuid ?? [])]) + pendingEvents.push({ + uuid, + changeType: response.changeType, + entity, + updatedFields: response.updatedFields, + }) + scheduleFlush() + } + }, + signal, + () => { connectionStatus = ConnectionStatus.DISCONNECTED } - } + ) } const streamSceneChanges = async (client: Client, signal: AbortSignal) => { - try { - for await (const response of client.streamSceneChanges({}, { signal })) { + await retryStream(async (sig) => { + for await (const response of client.streamSceneChanges({}, { signal: sig })) { const { sceneMetadata } = response if (!sceneMetadata) continue @@ -345,11 +347,7 @@ export function provideDrawService() { ) } } - } catch (error) { - if (!signal.aborted) { - console.error('Draw service scene stream error:', error) - } - } + }, signal) } $effect(() => { diff --git a/src/lib/retry-stream.ts b/src/lib/retry-stream.ts new file mode 100644 index 000000000..14a6e900f --- /dev/null +++ b/src/lib/retry-stream.ts @@ -0,0 +1,51 @@ +const INITIAL_DELAY_MS = 1_000 +const MAX_DELAY_MS = 30_000 + +/** + * Calls `run` in a loop, retrying with exponential backoff when it throws. + * - Clean stream end (server closed it): restarts immediately, delay resets. + * - Error: calls `onRetry`, waits with exponential backoff, then retries. + * Stops when the signal is aborted. + */ +export const retryStream = async ( + run: (signal: AbortSignal) => Promise, + signal: AbortSignal, + onRetry?: (delay: number) => void +): Promise => { + let delay = INITIAL_DELAY_MS + + while (!signal.aborted) { + let errored = false + try { + await run(signal) + // Stream ended cleanly (server closed it) — restart immediately. + delay = INITIAL_DELAY_MS + } catch (error) { + if (signal.aborted) return + errored = true + console.warn('Stream error, retrying in', delay, 'ms:', error) + } + + if (signal.aborted) return + + if (errored) { + onRetry?.(delay) + await sleep(delay, signal) + delay = Math.min(delay * 2, MAX_DELAY_MS) + } + } +} + +const sleep = (ms: number, signal: AbortSignal): Promise => { + return new Promise((resolve) => { + const timer = setTimeout(resolve, ms) + signal.addEventListener( + 'abort', + () => { + clearTimeout(timer) + resolve() + }, + { once: true } + ) + }) +} From 7120c0ed966aa7ccd71550372b81df9afa654944 Mon Sep 17 00:00:00 2001 From: "claude[bot]" <41898282+claude[bot]@users.noreply.github.com> Date: Wed, 8 Jul 2026 21:26:33 +0000 Subject: [PATCH 2/4] fix: bump svelte-sdk to 1.2.3, clear stale entities on retry, fix sleep leak - Bump @viamrobotics/svelte-sdk to 1.2.3 and set resetQueriesOnDisconnect:false in ViamProvider so query data persists across disconnects in the dev app; hook-level preserve logic remains as a safe fallback for library consumers - Clear transformEntities/drawingEntities and destroy ECS entities in the onRetry callback so the server's ADDED re-bootstrap after reconnect can rebuild from scratch (fixes silent entity drop on stale map entries) - Set connectionStatus = CONNECTING in onRetry so the UI shows reconnect is in progress rather than staying DISCONNECTED indefinitely - Add onRetry warn log to streamSceneChanges so scene-stream failures are no longer silently swallowed - Fix sleep() memory leak: remove the abort listener when the timer fires normally so it doesn't accumulate dangling listeners during long outages Co-authored-by: Devin T. Currie --- package.json | 2 +- pnpm-lock.yaml | 39 +++++++++---- .../DrawService/useDrawService.svelte.ts | 57 +++++++++++++------ src/lib/retry-stream.ts | 18 +++--- src/routes/+layout.svelte | 1 + 5 files changed, 79 insertions(+), 38 deletions(-) diff --git a/package.json b/package.json index 2a18cca60..dc62a7c34 100644 --- a/package.json +++ b/package.json @@ -485,7 +485,7 @@ "@typescript-eslint/parser": "8.56.1", "@viamrobotics/prime-core": "0.1.5", "@viamrobotics/sdk": "0.69.0", - "@viamrobotics/svelte-sdk": "1.2.2", + "@viamrobotics/svelte-sdk": "1.2.3", "@viamrobotics/tweakpane-config": "0.1.1", "@vitest/browser-playwright": "4.1.9", "@vitest/coverage-v8": "4.1.9", diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index f6688ca6d..b5ac40f40 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -161,8 +161,8 @@ importers: specifier: 0.69.0 version: 0.69.0 '@viamrobotics/svelte-sdk': - specifier: 1.2.2 - version: 1.2.2(@viamrobotics/sdk@0.69.0)(svelte@5.55.7) + specifier: 1.2.3 + version: 1.2.3(@sveltejs/kit@2.67.0(@opentelemetry/api@1.9.0)(@sveltejs/vite-plugin-svelte@7.1.2(svelte@5.55.7)(vite@8.1.0(@types/node@25.6.0)(esbuild@0.28.1)(jiti@2.7.0)(tsx@4.20.5)(yaml@2.8.1)))(svelte@5.55.7)(typescript@5.9.2)(vite@8.1.0(@types/node@25.6.0)(esbuild@0.28.1)(jiti@2.7.0)(tsx@4.20.5)(yaml@2.8.1)))(@viamrobotics/sdk@0.69.0)(svelte@5.55.7)(zod@4.4.3) '@viamrobotics/tweakpane-config': specifier: 0.1.1 version: 0.1.1(svelte-tweakpane-ui@1.5.16(svelte@5.55.7))(svelte@5.55.7) @@ -1995,8 +1995,8 @@ packages: '@viamrobotics/sdk@0.69.0': resolution: {integrity: sha512-xyJqF+4kJEmLiGiwWjf7hdcq/FBMcCcPF4kGyz3WQ/XeBC4HrkizD05ka6eNkevMFoAseGbNgOqOochEs0XTLA==} - '@viamrobotics/svelte-sdk@1.2.2': - resolution: {integrity: sha512-vuvmA+NZZcLPZs39o3vg4wlqaGF+xJU+lGGtKvwU/HJXpxuARVpPskTKvcjxN5RBjt9dQhjvQV/4er1tAlTAZA==} + '@viamrobotics/svelte-sdk@1.2.3': + resolution: {integrity: sha512-V7ejNb9fE8n5PRbdsDpyaycHIzYT/A80O73KoPNUk/XBRw7X4v1MRaWZuGuGjxOZpq/kwLo7eY6rqahWODHCOw==} peerDependencies: '@viamrobotics/sdk': '>=0.57' svelte: '>=5' @@ -3815,15 +3815,22 @@ packages: run-parallel@1.2.0: resolution: {integrity: sha512-5l4VyZR86LZ/lDxZTR6jqL8AFE2S0IFLMP26AbjsLVADxHdhB/c0GUsH+y39UfCi3dzz8OlQuPmnaJOMoDHQBA==} - runed@0.29.2: - resolution: {integrity: sha512-0cq6cA6sYGZwl/FvVqjx9YN+1xEBu9sDDyuWdDW1yWX7JF2wmvmVKfH+hVCZs+csW+P3ARH92MjI3H9QTagOQA==} + runed@0.31.1: + resolution: {integrity: sha512-v3czcTnO+EJjiPvD4dwIqfTdHLZ8oH0zJheKqAHh9QMViY7Qb29UlAMRpX7ZtHh7AFqV60KmfxaJ9QMy+L1igQ==} peerDependencies: svelte: ^5.7.0 - runed@0.31.1: - resolution: {integrity: sha512-v3czcTnO+EJjiPvD4dwIqfTdHLZ8oH0zJheKqAHh9QMViY7Qb29UlAMRpX7ZtHh7AFqV60KmfxaJ9QMy+L1igQ==} + runed@0.37.1: + resolution: {integrity: sha512-MeFY73xBW8IueWBm012nNFIGy19WUGPLtknavyUPMpnyt350M47PhGSGrGoSLbidwn+Zlt/O0cp8/OZE3LASWA==} peerDependencies: + '@sveltejs/kit': ^2.21.0 svelte: ^5.7.0 + zod: ^4.1.0 + peerDependenciesMeta: + '@sveltejs/kit': + optional: true + zod: + optional: true rxjs@7.8.2: resolution: {integrity: sha512-dhKf903U/PQZY6boNNtAGdWbG85WAbjT/1xYoZIC7FAY0yWapOBQVsVrDl58W86//e1VpMNBtRV4MaXfdMySFA==} @@ -6193,14 +6200,17 @@ snapshots: bsonfy: 1.0.2 exponential-backoff: 3.1.3 - '@viamrobotics/svelte-sdk@1.2.2(@viamrobotics/sdk@0.69.0)(svelte@5.55.7)': + '@viamrobotics/svelte-sdk@1.2.3(@sveltejs/kit@2.67.0(@opentelemetry/api@1.9.0)(@sveltejs/vite-plugin-svelte@7.1.2(svelte@5.55.7)(vite@8.1.0(@types/node@25.6.0)(esbuild@0.28.1)(jiti@2.7.0)(tsx@4.20.5)(yaml@2.8.1)))(svelte@5.55.7)(typescript@5.9.2)(vite@8.1.0(@types/node@25.6.0)(esbuild@0.28.1)(jiti@2.7.0)(tsx@4.20.5)(yaml@2.8.1)))(@viamrobotics/sdk@0.69.0)(svelte@5.55.7)(zod@4.4.3)': dependencies: '@tanstack/svelte-query': 6.1.28(svelte@5.55.7) '@tanstack/svelte-query-devtools': 6.1.28(@tanstack/svelte-query@6.1.28(svelte@5.55.7))(svelte@5.55.7) '@viamrobotics/sdk': 0.69.0 loglayer: 9.1.0 - runed: 0.29.2(svelte@5.55.7) + runed: 0.37.1(@sveltejs/kit@2.67.0(@opentelemetry/api@1.9.0)(@sveltejs/vite-plugin-svelte@7.1.2(svelte@5.55.7)(vite@8.1.0(@types/node@25.6.0)(esbuild@0.28.1)(jiti@2.7.0)(tsx@4.20.5)(yaml@2.8.1)))(svelte@5.55.7)(typescript@5.9.2)(vite@8.1.0(@types/node@25.6.0)(esbuild@0.28.1)(jiti@2.7.0)(tsx@4.20.5)(yaml@2.8.1)))(svelte@5.55.7)(zod@4.4.3) svelte: 5.55.7 + transitivePeerDependencies: + - '@sveltejs/kit' + - zod '@viamrobotics/tweakpane-config@0.1.1(svelte-tweakpane-ui@1.5.16(svelte@5.55.7))(svelte@5.55.7)': dependencies: @@ -8068,15 +8078,20 @@ snapshots: dependencies: queue-microtask: 1.2.3 - runed@0.29.2(svelte@5.55.7): + runed@0.31.1(svelte@5.55.7): dependencies: esm-env: 1.2.2 svelte: 5.55.7 - runed@0.31.1(svelte@5.55.7): + runed@0.37.1(@sveltejs/kit@2.67.0(@opentelemetry/api@1.9.0)(@sveltejs/vite-plugin-svelte@7.1.2(svelte@5.55.7)(vite@8.1.0(@types/node@25.6.0)(esbuild@0.28.1)(jiti@2.7.0)(tsx@4.20.5)(yaml@2.8.1)))(svelte@5.55.7)(typescript@5.9.2)(vite@8.1.0(@types/node@25.6.0)(esbuild@0.28.1)(jiti@2.7.0)(tsx@4.20.5)(yaml@2.8.1)))(svelte@5.55.7)(zod@4.4.3): dependencies: + dequal: 2.0.3 esm-env: 1.2.2 + lz-string: 1.5.0 svelte: 5.55.7 + optionalDependencies: + '@sveltejs/kit': 2.67.0(@opentelemetry/api@1.9.0)(@sveltejs/vite-plugin-svelte@7.1.2(svelte@5.55.7)(vite@8.1.0(@types/node@25.6.0)(esbuild@0.28.1)(jiti@2.7.0)(tsx@4.20.5)(yaml@2.8.1)))(svelte@5.55.7)(typescript@5.9.2)(vite@8.1.0(@types/node@25.6.0)(esbuild@0.28.1)(jiti@2.7.0)(tsx@4.20.5)(yaml@2.8.1)) + zod: 4.4.3 rxjs@7.8.2: dependencies: diff --git a/src/lib/plugins/DrawService/useDrawService.svelte.ts b/src/lib/plugins/DrawService/useDrawService.svelte.ts index 9aa6b2e7c..17983f59b 100644 --- a/src/lib/plugins/DrawService/useDrawService.svelte.ts +++ b/src/lib/plugins/DrawService/useDrawService.svelte.ts @@ -304,6 +304,19 @@ export function provideDrawService() { }) } + const clearEntities = () => { + for (const entity of transformEntities.values()) { + if (world.has(entity)) hierarchy.destroyEntityTree(world, entity) + } + transformEntities.clear() + + for (const entity of drawingEntities.values()) { + if (world.has(entity)) hierarchy.destroyEntityTree(world, entity) + } + drawingEntities.clear() + serverRelationships.reset() + } + const streamEntityChanges = async (client: Client, signal: AbortSignal) => { await retryStream( async (sig) => { @@ -325,29 +338,41 @@ export function provideDrawService() { }, signal, () => { - connectionStatus = ConnectionStatus.DISCONNECTED + // Set CONNECTING so the UI shows reconnect is in progress. + // Clear all entities so the server's re-bootstrap (ADDED for every + // live entity) can rebuild the scene from scratch. Without this, + // processTransformEvent/processDrawingEvent silently drop every ADDED + // for a UUID that's still in the maps, leaving the client permanently stale. + connectionStatus = ConnectionStatus.CONNECTING + clearEntities() } ) } const streamSceneChanges = async (client: Client, signal: AbortSignal) => { - await retryStream(async (sig) => { - for await (const response of client.streamSceneChanges({}, { signal: sig })) { - const { sceneMetadata } = response - if (!sceneMetadata) continue - - if (sceneMetadata.sceneCamera?.position && sceneMetadata.sceneCamera?.lookAt) { - const { position, lookAt, animated } = sceneMetadata.sceneCamera - cameraControls.setPose( - { - position: [position.x * 0.001, position.y * 0.001, position.z * 0.001], - lookAt: [lookAt.x * 0.001, lookAt.y * 0.001, lookAt.z * 0.001], - }, - animated ?? false - ) + await retryStream( + async (sig) => { + for await (const response of client.streamSceneChanges({}, { signal: sig })) { + const { sceneMetadata } = response + if (!sceneMetadata) continue + + if (sceneMetadata.sceneCamera?.position && sceneMetadata.sceneCamera?.lookAt) { + const { position, lookAt, animated } = sceneMetadata.sceneCamera + cameraControls.setPose( + { + position: [position.x * 0.001, position.y * 0.001, position.z * 0.001], + lookAt: [lookAt.x * 0.001, lookAt.y * 0.001, lookAt.z * 0.001], + }, + animated ?? false + ) + } } + }, + signal, + () => { + console.warn('Scene stream disconnected, retrying...') } - }, signal) + ) } $effect(() => { diff --git a/src/lib/retry-stream.ts b/src/lib/retry-stream.ts index 14a6e900f..e5908754f 100644 --- a/src/lib/retry-stream.ts +++ b/src/lib/retry-stream.ts @@ -38,14 +38,14 @@ export const retryStream = async ( const sleep = (ms: number, signal: AbortSignal): Promise => { return new Promise((resolve) => { - const timer = setTimeout(resolve, ms) - signal.addEventListener( - 'abort', - () => { - clearTimeout(timer) - resolve() - }, - { once: true } - ) + const onAbort = () => { + clearTimeout(timer) + resolve() + } + const timer = setTimeout(() => { + signal.removeEventListener('abort', onAbort) + resolve() + }, ms) + signal.addEventListener('abort', onAbort, { once: true }) }) } diff --git a/src/routes/+layout.svelte b/src/routes/+layout.svelte index 49132423d..9eecc5889 100644 --- a/src/routes/+layout.svelte +++ b/src/routes/+layout.svelte @@ -51,6 +51,7 @@ }, }, }} + options={{ resetQueriesOnDisconnect: false }} {dialConfigs} > Date: Wed, 15 Jul 2026 18:54:56 +0000 Subject: [PATCH 3/4] test: clarify misleading retry-stream test name Rename 'stops retrying when signal is aborted' to accurately reflect that onRetry IS called once (signaling intent to retry) before the backoff sleep, but run is not called again because the abort fires during the backoff window. Co-authored-by: Devin T. Currie --- src/lib/__tests__/retry-stream.spec.ts | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/lib/__tests__/retry-stream.spec.ts b/src/lib/__tests__/retry-stream.spec.ts index dcaa7347c..568dcfccd 100644 --- a/src/lib/__tests__/retry-stream.spec.ts +++ b/src/lib/__tests__/retry-stream.spec.ts @@ -43,7 +43,7 @@ describe('retryStream', () => { vi.useRealTimers() }) - it('stops retrying when signal is aborted', async () => { + it('calls onRetry but does not retry run when signal is aborted during backoff', async () => { vi.useFakeTimers() const controller = new AbortController() From 4e925fb16f5e8a16b4bebf340d1c1c2e8468ae6b Mon Sep 17 00:00:00 2001 From: "claude[bot]" <41898282+claude[bot]@users.noreply.github.com> Date: Fri, 17 Jul 2026 18:55:49 +0000 Subject: [PATCH 4/4] test: add MAX_DELAY_MS cap and pre-aborted signal tests for retryStream Co-authored-by: Devin T. Currie --- src/lib/__tests__/retry-stream.spec.ts | 42 ++++++++++++++++++++++++++ 1 file changed, 42 insertions(+) diff --git a/src/lib/__tests__/retry-stream.spec.ts b/src/lib/__tests__/retry-stream.spec.ts index 568dcfccd..d1c295848 100644 --- a/src/lib/__tests__/retry-stream.spec.ts +++ b/src/lib/__tests__/retry-stream.spec.ts @@ -115,6 +115,48 @@ describe('retryStream', () => { expect(onRetry).not.toHaveBeenCalled() }) + it('caps delay at MAX_DELAY_MS (30s) and does not double beyond it', async () => { + vi.useFakeTimers() + + const controller = new AbortController() + let callCount = 0 + + const run = vi.fn().mockImplementation(async () => { + callCount++ + if (callCount <= 6) throw new Error('stream error') + controller.abort() + }) + + const onRetry = vi.fn() + const promise = retryStream(run, controller.signal, onRetry) + + // Advance through each exponential backoff step + await vi.advanceTimersByTimeAsync(1_000) // delay: 1000 + await vi.advanceTimersByTimeAsync(2_000) // delay: 2000 + await vi.advanceTimersByTimeAsync(4_000) // delay: 4000 + await vi.advanceTimersByTimeAsync(8_000) // delay: 8000 + await vi.advanceTimersByTimeAsync(16_000) // delay: 16000 + await vi.advanceTimersByTimeAsync(30_000) // delay: capped at 30000 (not 32000) + + await promise + + expect(onRetry).toHaveBeenCalledTimes(6) + expect(onRetry).toHaveBeenNthCalledWith(5, 16_000) + expect(onRetry).toHaveBeenNthCalledWith(6, 30_000) + + vi.useRealTimers() + }) + + it('does not call run when signal is already aborted before retryStream is called', async () => { + const controller = new AbortController() + controller.abort() + + const run = vi.fn() + await retryStream(run, controller.signal) + + expect(run).not.toHaveBeenCalled() + }) + it('resets delay after a successful run', async () => { vi.useFakeTimers()