From 418b246b34410bdc97d5e7aac66fd30311e8a4a3 Mon Sep 17 00:00:00 2001 From: joaner Date: Tue, 11 Aug 2026 09:53:15 +0800 Subject: [PATCH] fix(bag): stop infinite Range re-fetch loop when loading remote .bag files MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit @foxglove/rosbag's Filelike.size() must be synchronous (number), but the remote adapter in bag.worker.ts exposed an async/BigInt size(), which turned the library's un-awaited `size() - offset` arithmetic into NaN. A NaN-bounded read() can never be marked satisfied by CachedFilelike, so it kept re-deriving and re-fetching the same ~50MiB block forever while Bag.open() never resolved — reproduced end-to-end against a real 480MB bag file (initialize() hung past 12s and even emitted a literal `Range: bytes=NaN-NaN` request before the fix; resolves in ~500ms after). Also hardens CachedFilelike against this whole class of bug: reject non-finite/negative offsets and lengths at the read()/prefetch() entry points instead of silently enqueueing unsatisfiable requests, and bound the connection-error retry count so a persistently (but not rapidly) failing fetch can no longer retry forever. Bump to v1.7.9. --- package-lock.json | 4 +- package.json | 2 +- src/infra/services/CachedFilelike.test.ts | 92 ++++++++++++++++++++ src/infra/services/CachedFilelike.ts | 96 +++++++++++++++++---- src/infra/sources/BagIterableSource.ts | 22 +++-- src/infra/workers/bag.worker.ts | 10 +-- src/infra/workers/remoteBagReadable.test.ts | 60 +++++++++++++ src/infra/workers/remoteBagReadable.ts | 32 +++++++ 8 files changed, 280 insertions(+), 38 deletions(-) create mode 100644 src/infra/workers/remoteBagReadable.test.ts create mode 100644 src/infra/workers/remoteBagReadable.ts diff --git a/package-lock.json b/package-lock.json index b432a38..234c55a 100644 --- a/package-lock.json +++ b/package-lock.json @@ -1,12 +1,12 @@ { "name": "@ioai/rosview", - "version": "1.7.8", + "version": "1.7.9", "lockfileVersion": 3, "requires": true, "packages": { "": { "name": "@ioai/rosview", - "version": "1.7.8", + "version": "1.7.9", "license": "MIT", "devDependencies": { "@eslint/js": "^10.0.1", diff --git a/package.json b/package.json index 470b3ce..b758ab5 100644 --- a/package.json +++ b/package.json @@ -1,6 +1,6 @@ { "name": "@ioai/rosview", - "version": "1.7.8", + "version": "1.7.9", "description": "High-performance robotics data visualization for MCAP, ROS bag, ROS2 db3, HDF5 and BVH — embeddable React component and standalone SPA", "keywords": [ "ros", diff --git a/src/infra/services/CachedFilelike.test.ts b/src/infra/services/CachedFilelike.test.ts index 45a5a22..e9b58ab 100644 --- a/src/infra/services/CachedFilelike.test.ts +++ b/src/infra/services/CachedFilelike.test.ts @@ -134,3 +134,95 @@ describe('CachedFilelike prefetch', () => { await expect(readPromise).resolves.toEqual(new Uint8Array([3])); }); }); + +// Regression tests for a bug where a remote `Filelike` adapter with a mismatched (async/bigint) +// `size()` caused `read(offset, NaN)` calls to hang forever while `CachedFilelike` kept +// re-fetching an already-downloaded ~50MiB block on a loop, with no error ever surfacing. See +// `remoteBagReadable.test.ts` for the corresponding regression test at the adapter boundary. +describe('CachedFilelike input validation', () => { + it('rejects non-finite read lengths synchronously instead of enqueueing an unsatisfiable request', async () => { + const reader = new TestFileReader(); + const filelike = new CachedFilelike({ fileReader: reader, cacheSizeInBytes: 32, fetchBlockSizeInBytes: 8 }); + + expect(() => filelike.read(4, NaN)).toThrow(/invalid input/); + expect(() => filelike.read(NaN, 4)).toThrow(/invalid input/); + expect(() => filelike.read(0, Infinity)).toThrow(/invalid input/); + await flushAsyncWork(); + + // No fetch should ever be scheduled for an unsatisfiable range. + expect(reader.streams).toHaveLength(0); + }); + + it('rejects negative or non-integer offsets/lengths', () => { + const reader = new TestFileReader(); + const filelike = new CachedFilelike({ fileReader: reader, cacheSizeInBytes: 32 }); + + expect(() => filelike.read(-1, 4)).toThrow(/invalid input/); + expect(() => filelike.read(0, -4)).toThrow(/invalid input/); + expect(() => filelike.read(1.5, 4)).toThrow(/invalid input/); + }); + + it('silently drops malformed prefetch requests instead of scheduling a fetch', async () => { + const reader = new TestFileReader(); + const filelike = new CachedFilelike({ fileReader: reader, cacheSizeInBytes: 32, fetchBlockSizeInBytes: 8 }); + + filelike.prefetch(4, NaN); + filelike.prefetch(-1, 4); + filelike.prefetch(1.5, 4); + await flushAsyncWork(); + + expect(reader.streams).toHaveLength(0); + }); +}); + +describe('CachedFilelike avoids re-fetching already-satisfied ranges', () => { + it('does not issue a second fetch for a block that is already fully downloaded', async () => { + const reader = new TestFileReader(); + const filelike = new CachedFilelike({ fileReader: reader, cacheSizeInBytes: 32, fetchBlockSizeInBytes: 8 }); + + const firstRead = filelike.read(0, 8); + await flushAsyncWork(); + expect(reader.streams).toHaveLength(1); + + reader.streams[0].emitData([0, 1, 2, 3, 4, 5, 6, 7]); + await expect(firstRead).resolves.toEqual(new Uint8Array([0, 1, 2, 3, 4, 5, 6, 7])); + + // Requesting the exact same already-downloaded range again must be served from cache, + // not trigger a new HTTP-equivalent fetch. + const secondRead = await filelike.read(0, 8); + expect(secondRead).toEqual(new Uint8Array([0, 1, 2, 3, 4, 5, 6, 7])); + expect(reader.streams).toHaveLength(1); + }); +}); + +describe('CachedFilelike bounded error retries', () => { + it('gives up after a bounded number of consecutive errors, even when failures are spaced beyond the 100ms rapid-fault window', async () => { + let now = 0; + const dateNowSpy = vi.spyOn(Date, 'now').mockImplementation(() => now); + try { + const reader = new TestFileReader(); + const filelike = new CachedFilelike({ fileReader: reader, cacheSizeInBytes: 32, fetchBlockSizeInBytes: 8 }); + + const readPromise = filelike.read(0, 8); + await flushAsyncWork(); + + let settled = false; + void readPromise.catch(() => { + settled = true; + }); + + for (let i = 0; i < 20 && !settled; i++) { + const stream = reader.streams[reader.streams.length - 1]; + now += 200; // well beyond the old 100ms rapid-double-fault window + stream.emit('error', new Error(`boom ${i}`)); + await flushAsyncWork(); + } + + await expect(readPromise).rejects.toThrow(/giving up/); + // The retry budget must be small and bounded — not "keep retrying forever". + expect(reader.streams.length).toBeLessThan(20); + } finally { + dateNowSpy.mockRestore(); + } + }); +}); diff --git a/src/infra/services/CachedFilelike.ts b/src/infra/services/CachedFilelike.ts index 38c6398..822a8fe 100644 --- a/src/infra/services/CachedFilelike.ts +++ b/src/infra/services/CachedFilelike.ts @@ -20,6 +20,19 @@ export interface FileReader { const CACHE_BLOCK_SIZE = 1024 * 1024 * 50; // 50MiB blocks const DEFAULT_MAX_REQUEST_SIZE = CACHE_BLOCK_SIZE * 2; +/** + * Max consecutive fetch failures for the *same* logical block before giving up and rejecting + * pending reads. Previously the only give-up condition was "two errors within 100ms of each + * other", which never triggers against a server/network that fails slowly-but-persistently + * (RTT > 100ms) — that failure mode retried forever. This bound guarantees termination + * regardless of error timing. + */ +const MAX_CONSECUTIVE_BLOCK_ERRORS = 6; + +function isFiniteNonNegativeInteger(value: number): boolean { + return Number.isInteger(value) && value >= 0; +} + export default class CachedFilelike implements Readable { #fileReader: FileReader; #cacheSizeInBytes: number = Infinity; @@ -42,6 +55,7 @@ export default class CachedFilelike implements Readable { #prefetchRequests: Range[] = []; #lastErrorTime?: number; + #consecutiveBlockErrorCount: number = 0; public constructor(options: { fileReader: FileReader; @@ -107,11 +121,23 @@ export default class CachedFilelike implements Readable { return Promise.resolve(new Uint8Array()); } + // Fail fast on non-finite / negative / non-integer offsets or lengths (e.g. `NaN` from a + // caller doing arithmetic on an un-awaited `Promise`). Without this guard, a `NaN` `end` + // silently turns into a `read()` that can never be satisfied — `hasData()` is never true + // for a `NaN` bound — while the block-alignment logic keeps computing a plausible-looking, + // finite fetch range from the (valid) `start` and re-requesting it forever. See + // `bag.worker.ts`'s remote `Filelike.size()` adapter for the real-world case this fixes. + if ( + !isFiniteNonNegativeInteger(offset) || + !isFiniteNonNegativeInteger(length) + ) { + throw new Error( + `CachedFilelike#read invalid input: offset=${offset}, length=${length} (must be finite non-negative integers)`, + ); + } + const range = { start: offset, end: offset + length }; - if (offset < 0 || length < 0) { - throw new Error("CachedFilelike#read invalid input"); - } if (length > this.#cacheSizeInBytes) { throw new Error(`Requested more data than cache size: ${length} > ${this.#cacheSizeInBytes}`); } @@ -138,12 +164,18 @@ export default class CachedFilelike implements Readable { if (length <= 0 || this.#closed) { return; } - - const range = { start: offset, end: offset + length }; - if (offset < 0 || length < 0 || length > this.#cacheSizeInBytes) { + // Best-effort: silently drop malformed prefetch requests rather than let a `NaN`/negative + // bound reach the same range-alignment code path that `read()` guards against above. + if ( + !isFiniteNonNegativeInteger(offset) || + !isFiniteNonNegativeInteger(length) || + length > this.#cacheSizeInBytes + ) { return; } + const range = { start: offset, end: offset + length }; + void this.open() .then(async () => { const size = await this.size(); @@ -177,7 +209,14 @@ export default class CachedFilelike implements Readable { return; } - this.#readRequests = this.#readRequests.filter(({ range, resolve }) => { + this.#readRequests = this.#readRequests.filter(({ range, resolve, reject }) => { + // Second line of defense: `read()` already rejects non-finite ranges synchronously, but + // reject here too in case a request ever reaches the queue some other way — an + // unsatisfiable range must never sit in the queue silently forever. + if (!Number.isFinite(range.start) || !Number.isFinite(range.end)) { + reject(new Error(`CachedFilelike: unsatisfiable range [${range.start}, ${range.end})`)); + return false; + } if (!this.#virtualBuffer.hasData(range.start, range.end)) { return true; } @@ -227,6 +266,12 @@ export default class CachedFilelike implements Readable { } #getNextFixedFetchRange(queryRange: Range, fileSize: number): Range | undefined { + if (!Number.isFinite(queryRange.start) || !Number.isFinite(queryRange.end)) { + // Never schedule a fetch from a non-finite range: `#alignToFetchBlock` would otherwise + // happily derive a plausible, finite block from `queryRange.start` alone and re-fetch it + // forever, since the (bogus) query range itself can never be marked satisfied. + return undefined; + } if (queryRange.start >= fileSize) { return undefined; } @@ -289,19 +334,31 @@ export default class CachedFilelike implements Readable { return; } - if (this.#keepReconnectingCallback) { - if (this.#lastErrorTime == undefined) { - this.#keepReconnectingCallback(true); - } - } else { - const lastErrorTime = this.#lastErrorTime; - if (lastErrorTime != undefined && Date.now() - lastErrorTime < 100) { - this.#closed = true; - for (const request of this.#readRequests) { - request.reject(error); - } - return; + // Bounded regardless of `keepReconnectingCallback` and independent of the "two errors + // within 100ms" heuristic below, which never trips against a slowly-but-persistently + // failing server/network (RTT > 100ms) — that combination used to retry forever. + this.#consecutiveBlockErrorCount += 1; + const exhaustedRetryBudget = this.#consecutiveBlockErrorCount >= MAX_CONSECUTIVE_BLOCK_ERRORS; + const rapidDoubleFault = + !this.#keepReconnectingCallback && + this.#lastErrorTime != undefined && + Date.now() - this.#lastErrorTime < 100; + + if (exhaustedRetryBudget || rapidDoubleFault) { + this.#closed = true; + const failure = exhaustedRetryBudget + ? new Error( + `CachedFilelike: giving up on ${range.start}-${range.end} after ${this.#consecutiveBlockErrorCount} consecutive errors: ${error.message}`, + ) + : error; + for (const request of this.#readRequests) { + request.reject(failure); } + return; + } + + if (this.#keepReconnectingCallback && this.#lastErrorTime == undefined) { + this.#keepReconnectingCallback(true); } this.#lastErrorTime = Date.now(); @@ -317,6 +374,7 @@ export default class CachedFilelike implements Readable { return; } + this.#consecutiveBlockErrorCount = 0; if (this.#lastErrorTime != undefined) { this.#lastErrorTime = undefined; if (this.#keepReconnectingCallback) { diff --git a/src/infra/sources/BagIterableSource.ts b/src/infra/sources/BagIterableSource.ts index e4bf70b..a66fd15 100644 --- a/src/infra/sources/BagIterableSource.ts +++ b/src/infra/sources/BagIterableSource.ts @@ -1,4 +1,4 @@ -import { Bag } from "@foxglove/rosbag"; +import { Bag, type Filelike } from "@foxglove/rosbag"; import { BlobReader } from "@foxglove/rosbag/web"; import { parse as parseMessageDefinition } from "@foxglove/rosmsg"; import { MessageReader } from "@foxglove/rosmsg-serialization"; @@ -15,11 +15,14 @@ import type { MessageIteratorArgs, GetBackfillMessagesArgs } from '@/infra/worke import { loadDecompressHandlers } from "./decompressHandlers"; import { addMs, toNano } from '@/shared/utils/time'; -/** Remote byte reader shape accepted by @foxglove/rosbag `Bag` (non-`BlobReader` paths). */ -interface RemoteBagReadable { - size: () => Promise; - read: (offset: number, length: number) => Promise; -} +/** + * Remote byte reader shape accepted by @foxglove/rosbag `Bag` (non-`BlobReader` paths). + * Must structurally match the library's own `Filelike` (`size(): number`, synchronous) — + * previously this was hand-rolled with an async/bigint `size()`, which the compiler could + * not catch because callers cast past it (see the removed `as ConstructorParameters<...>` + * below). Deriving from `Filelike` directly means any future mismatch is a type error again. + */ +type RemoteBagReadable = Filelike; type BagSource = { type: "file"; file: Blob } | { type: "remote"; readable: RemoteBagReadable }; @@ -60,11 +63,12 @@ export class BagIterableSource implements IIterableSource { async initialize(): Promise { const decompressHandlers = await loadDecompressHandlers({ wasmBinary: this._wasmBinary }); - const fileLike: BlobReader | RemoteBagReadable = + const fileLike: Filelike = this._source.type === "remote" ? this._source.readable : new BlobReader(this._source.file); - // Rosbag `Bag` accepts `BlobReader` or custom readers; remote readers use bigint `size()` which differs from `BlobReader` typing. - this._bag = new Bag(fileLike as ConstructorParameters[0], { + // `BlobReader` and `RemoteBagReadable` both satisfy `Filelike` structurally now, so this + // no longer needs an `as`-cast to bypass the type checker. + this._bag = new Bag(fileLike, { parse: false, decompress: { // RosView currently supports lz4-compressed ROS1 bag chunks. diff --git a/src/infra/workers/bag.worker.ts b/src/infra/workers/bag.worker.ts index f7daf73..955b528 100644 --- a/src/infra/workers/bag.worker.ts +++ b/src/infra/workers/bag.worker.ts @@ -19,6 +19,7 @@ import type { TransportDiagnostics, WorkerTransportConfig } from "./transport"; import { SharedPayloadRing } from "./sharedPayloadRing"; import { resolveRemoteCacheBytes } from './remoteCacheConfig'; import { DataQualityScanController } from './dataQualityScanController'; +import { buildRemoteBagReadable, type SyncSizeBagReadable } from './remoteBagReadable'; class BagWorker implements IWorkerSerializedSourceWorker { private _source?: BagIterableSource; @@ -35,7 +36,7 @@ class BagWorker implements IWorkerSerializedSourceWorker { const url = typeof args.url === 'string' ? args.url : undefined; const file = args.file instanceof Blob ? args.file : undefined; let sourceArgs: - | { type: 'remote'; readable: { size: () => Promise; read: (offset: number, length: number) => Promise } } + | { type: 'remote'; readable: SyncSizeBagReadable } | { type: 'file'; file: Blob }; if (url) { const knownRaw = args.knownTotalBytes; @@ -54,12 +55,7 @@ class BagWorker implements IWorkerSerializedSourceWorker { cacheSizeInBytes: resolveRemoteCacheBytes(), }); this._cachedReadable = readable; - // We need to implement Filelike interface for rosbag - // For now, wrap it in an object that rosbag expects - const bagReadable = { - size: async () => BigInt(await readable.size()), - read: async (offset: number, length: number) => await readable.read(offset, length) - }; + const bagReadable = await buildRemoteBagReadable(readable); sourceArgs = { type: "remote", readable: bagReadable }; } else if (file) { this._cachedReadable = undefined; diff --git a/src/infra/workers/remoteBagReadable.test.ts b/src/infra/workers/remoteBagReadable.test.ts new file mode 100644 index 0000000..99768ab --- /dev/null +++ b/src/infra/workers/remoteBagReadable.test.ts @@ -0,0 +1,60 @@ +import { describe, expect, it } from 'vitest'; + +import { buildRemoteBagReadable, type AsyncByteRangeSource } from './remoteBagReadable'; + +/** Minimal `CachedFilelike`-shaped stub: async `size()`, async `read()`. */ +class TestAsyncByteRangeSource implements AsyncByteRangeSource { + constructor(private readonly totalBytes: number) {} + + async size(): Promise { + return this.totalBytes; + } + + async read(offset: number, length: number): Promise { + return new Uint8Array(length).fill(offset % 256); + } +} + +describe('buildRemoteBagReadable', () => { + it('exposes size() synchronously as a plain number, matching @foxglove/rosbag Filelike', async () => { + const readable = await buildRemoteBagReadable(new TestAsyncByteRangeSource(1234)); + + // `Filelike.size()` must be callable without `await` and must not be a `Promise` or + // `bigint` — @foxglove/rosbag calls it directly, e.g. `this._file.size() - fileOffset`. + // Regression test for the bug where `size` was `async () => BigInt(await source.size())`: + // that turned every un-awaited `size() - offset` into `NaN`, which made `CachedFilelike` + // re-fetch the same block forever because a `NaN`-bounded read can never be satisfied. + const size = readable.size(); + expect(typeof size).toBe('number'); + expect(size).toBe(1234); + + // The exact arithmetic @foxglove/rosbag performs when reading the connections/chunk + // index (`this._file.size() - fileOffset`) must be a finite number, never `NaN`. + const fileOffset = 100; + expect(Number.isFinite(readable.size() - fileOffset)).toBe(true); + expect(readable.size() - fileOffset).toBe(1134); + }); + + it('only resolves the underlying async size() once, then reuses it synchronously', async () => { + let sizeCalls = 0; + const source: AsyncByteRangeSource = { + size: async () => { + sizeCalls += 1; + return 42; + }, + read: async (_offset, length) => new Uint8Array(length), + }; + + const readable = await buildRemoteBagReadable(source); + expect(readable.size()).toBe(42); + expect(readable.size()).toBe(42); + expect(readable.size()).toBe(42); + expect(sizeCalls).toBe(1); + }); + + it('forwards read() to the underlying source', async () => { + const readable = await buildRemoteBagReadable(new TestAsyncByteRangeSource(10)); + const data = await readable.read(3, 4); + expect(data).toEqual(new Uint8Array([3, 3, 3, 3])); + }); +}); diff --git a/src/infra/workers/remoteBagReadable.ts b/src/infra/workers/remoteBagReadable.ts new file mode 100644 index 0000000..d533b40 --- /dev/null +++ b/src/infra/workers/remoteBagReadable.ts @@ -0,0 +1,32 @@ +/** Byte-range source with an async `size()` — matches `CachedFilelike`. */ +export interface AsyncByteRangeSource { + size(): Promise; + read(offset: number, length: number): Promise; +} + +/** `Filelike`-shaped adapter `@foxglove/rosbag`'s `Bag` requires for a remote source. */ +export interface SyncSizeBagReadable { + size: () => number; + read: (offset: number, length: number) => Promise; +} + +/** + * Builds the `Filelike`-shaped adapter that `@foxglove/rosbag`'s `Bag` requires for a remote + * source. `Filelike.size()` MUST be synchronous (`number`, never `Promise`/`bigint`) — the + * library calls it directly without `await`, e.g. `this._file.size() - fileOffset` when + * reading the connections/chunk index. Resolving the byte count once up front (rather than + * exposing an `async`/`Promise`-returning `size()`) is what keeps that arithmetic from + * silently becoming `NaN`. A `NaN` length there used to cause `CachedFilelike` to re-fetch the + * same ~50MiB block forever, since a `NaN`-bounded read can never be marked satisfied (see + * `CachedFilelike`'s input validation for the second line of defense against this class of + * bug). Kept in its own module (rather than inline in `bag.worker.ts`) so the sync contract + * can be regression-tested without importing a worker entrypoint that calls `Comlink.expose` + * at module load time. + */ +export async function buildRemoteBagReadable(source: AsyncByteRangeSource): Promise { + const totalBytes = await source.size(); + return { + size: () => totalBytes, + read: (offset: number, length: number) => source.read(offset, length), + }; +}