fix(puter-js): request idle timeout, callback/listener leaks, guarded message handlers (PUT-1908) (#4143)

* fix(puter-js): free one-shot callbacks and listeners, guard message handlers

- setAuthToken starts the GUI cache refresh loop once instead of per call
- IPC reply callbacks are dropped after their reply; dehydrated callbacks stay
- AppConnection removes its window listener when the target app closes
- reply branches for app data and instance count drop their callback
- UI, Debug, tool and GUI-token message listeners no longer reject; an
  unknown or failing tool answers toolResponse with an error
- parseResponse handles a blob response with no content-type
- retry waits remove their abort listener when the wait ends
- EmailConfirmationDialog escapes its message
- EventListener uses a real Map and snapshots listeners per emit
- Context get/set/runWithContext route every known key the same way

* fix(puter-js): time out requests that make no progress

- buildXhr starts an idle clock on send, reset by state changes and body
  progress (upload progress only on requests already carrying a bearer)
- read-safe requests get 60 s while progress is observable, everything
  else 15 min; timeout: 0 on fetchUrl/driverCall/sendWithRetry turns it off
- a timed-out read replays once; writes and AI calls never replay
- timeouts reject with code request_timeout on fetchUrl (TypeError),
  driverCall, streams and the legacy XHR path
- document the limits in rate-limits-and-quotas and the chat stream note
This commit is contained in:
Daniel Salazar authored and GitHub committed 2026-10-09 00:07:20 -07:00
1 parent ae234ffcc2
commit 8ecfeec8c3
26 files changed
+1447 -151

No files matched your search

+52
View File
@@ -0,0 +1,52 @@
import { describe, expect, it } from 'vitest';
import { Context, runInDerivedContext, runWithContext } from './context';
describe('Context', () => {
it('reads back a known key set after the scope opened with it', () => {
runWithContext({ driverName: 'initial' }, () => {
Context.set('driverName', 'addressed');
expect(Context.get('driverName')).toBe('addressed');
});
});
it('reads every known key from where set writes it', () => {
const signal = new AbortController().signal;
runWithContext({}, () => {
Context.set('requestId', 'r-1');
Context.set('driverName', 'd');
Context.set('abortSignal', signal);
Context.set('strictUpstreamErrors', true);
expect(Context.get('requestId')).toBe('r-1');
expect(Context.get('driverName')).toBe('d');
expect(Context.get('abortSignal')).toBe(signal);
expect(Context.get('strictUpstreamErrors')).toBe(true);
expect(Context.current()?.extra.size).toBe(0);
});
});
it('keeps ad-hoc keys in the open map', () => {
runWithContext({}, () => {
Context.set('myService.txId', 'tx');
expect(Context.get('myService.txId')).toBe('tx');
expect(Context.get('toString')).toBeUndefined();
});
});
it('keeps an untyped extra key passed to the scope readable', () => {
const initial = { requestId: 'r', 'myService.txId': 'tx' };
runWithContext(initial as { requestId: string }, () => {
expect(Context.get('requestId')).toBe('r');
expect(Context.get('myService.txId')).toBe('tx');
});
});
it('isolates a derived scope from its parent', () => {
runWithContext({ driverName: 'outer' }, () => {
runInDerivedContext(() => {
Context.set('driverName', 'inner');
expect(Context.get('driverName')).toBe('inner');
});
expect(Context.get('driverName')).toBe('outer');
});
});
});
+29 -12
View File
@@ -76,6 +76,18 @@ export interface KnownContextFields {
strictUpstreamErrors: boolean;
}
// Every key of `KnownContextFields`, checked both ways by the compiler, so
// `get` and `set` can never disagree about where a known key lives.
const KNOWN_KEY_FLAGS: Record<keyof KnownContextFields, true> = {
actor: true,
req: true,
requestId: true,
driverName: true,
abortSignal: true,
strictUpstreamErrors: true,
};
const KNOWN_KEYS: ReadonlySet<string> = new Set(Object.keys(KNOWN_KEY_FLAGS));
// -- Context store ---------------------------------------------------
interface ContextStore {
@@ -85,13 +97,21 @@ interface ContextStore {
const als = new AsyncLocalStorage<ContextStore>();
const storeValue = (store: ContextStore, key: string, value: unknown) => {
if (KNOWN_KEYS.has(key)) {
(store.known as Record<string, unknown>)[key] = value;
} else {
store.extra.set(key, value);
}
};
// -- Public API ------------------------------------------------------
/**
* Static-style context accessor.
*
* Well-known keys (`actor`, `req`, `requestId`) return typed values. Any other
* string key hits the generic map and returns `unknown`.
* Well-known keys (`KnownContextFields`) return typed values. Any other string
* key hits the generic map and returns `unknown`.
*/
export class Context {
/**
@@ -111,7 +131,7 @@ export class Context {
if (key === undefined) return als.getStore();
const store = als.getStore();
if (!store) return undefined;
if (key in store.known) {
if (KNOWN_KEYS.has(key)) {
return (store.known as Record<string, unknown>)[key];
}
return store.extra.get(key);
@@ -134,11 +154,7 @@ export class Context {
`Context.set('${key}', ...) called outside a request scope`,
);
}
if (key === 'actor' || key === 'req' || key === 'requestId') {
(store.known as Record<string, unknown>)[key] = value;
} else {
store.extra.set(key, value);
}
storeValue(store, key, value);
}
/**
@@ -174,9 +190,10 @@ export const runWithContext = <T>(
initial: Partial<KnownContextFields>,
fn: () => T,
): T => {
const store: ContextStore = {
known: { ...initial },
extra: new Map(),
};
const store: ContextStore = { known: {}, extra: new Map() };
// Routed like `set`, so an untyped caller's extra key stays readable.
for (const [key, value] of Object.entries(initial)) {
storeValue(store, key, value);
}
return als.run(store, fn);
};
+6 -2
View File
@@ -123,8 +123,12 @@ if the provider stops sending the response before it is complete. Check for
error chunks in your `for await...of` loop before treating the response as complete.
If the network connection fails after streaming starts, the loop throws
`{ message, code: "network_error" }`. Use `try...catch` around the call and loop
to handle these failures.
`{ message, code: "network_error" }`. If the stream sends nothing for 15
minutes, it throws `{ message, code: "request_timeout" }`; a call that gets no
response at all for 15 minutes rejects the same way. Neither is retried, since
the model may already have run (see
[Request timeouts](/rate-limits-and-quotas/#request-timeouts)). Use
`try...catch` around the call and loop to handle these failures.
### Stopping a stream
+14
View File
@@ -512,6 +512,19 @@ Every account has a filesystem quota (100 MiB free; paid plans add more). It cou
An upload reserves its declared size, minus the size of any file it replaces, from the moment it starts. The reservation is released when the upload completes (the real size counts instead), is cancelled, or expires (15 minutes after starting by default, at most 1 hour, plus 5 minutes' grace). An abandoned upload holds its space until it expires. A `startBatchWrite` that doesn't fit fails as a whole, up front. `space()` counts stored bytes only, not reservations.
## Request timeouts
Puter.js stops a request that makes no progress for too long and rejects it with `code: "request_timeout"`. The clock restarts whenever the request moves: a state change, response bytes arriving, or upload progress where the runtime reports it. A download or stream that keeps moving is never cut off, however long it takes. The file contents `puter.fs.upload()` and `puter.fs.write()` send are not subject to this limit.
| Request | Stopped after this long without progress |
| ----------------------------------------------------------------------------------------------------------------- | ---------------------------------------- |
| Reads: read-only calls such as `puter.kv.get()`, `puter.kv.list()`, `puter.apps.get()` and `puter.hosting.list()` | 60 s |
| Everything else, including writes, AI calls, streamed responses and every `puter.fs` call | 15 min |
A read is held to 60 s only while progress can be seen: before the response headers arrive, and after them once the body has started arriving. Node.js, service workers and Puter Workers report no progress for a non-streamed body, so there a read moves to the 15 min limit once its headers are in.
A read that times out is retried once. A write or AI call that times out is never retried, because the server may already have done the work: check before repeating it. A stream that stalls throws `{ message, code: "request_timeout" }` from its loop.
## What happens when you hit a limit
| Status | `code` | Meaning | What to do |
@@ -520,6 +533,7 @@ An upload reserves its declared size, minus the size of any file it replaces, fr
| `402` | `insufficient_funds` | Monthly credit spent | The user buys credit or upgrades; the allowance resets next month. |
| `402` | `subscription_required` | The endpoint requires a paid plan | The user upgrades. Retrying or waiting doesn't help. |
| `413` | `storage_limit_reached` | Storage quota reached | The user deletes files or upgrades. |
| none | `request_timeout` | No progress within the [request timeout](#request-timeouts) | Retry if the call is safe to repeat. A write or AI call may already have run. |
Errors are JSON: `{ "error": …, "message": …, "code": … }`.
+29 -42
View File
@@ -2,6 +2,7 @@ import kvjs from '@heyputer/kv.js';
import APICallLogger from './lib/APICallLogger.js';
import { hasOpaqueOrigin } from './lib/auth-popup.js';
import { fetchUrl } from './lib/networkUtils.js';
import { handleToolRequest } from './lib/toolBridge.js';
import { isStoredTokenUsableForOrigin } from './lib/authTokenOrigin.js';
import { isFramedDocument } from './lib/appModeGate.js';
import {
@@ -1034,10 +1035,13 @@ export class Puter {
console.error('Error accessing localStorage:', error);
}
}
// initialize loop for updating caches for major directories
if (this.env === 'gui') {
// check and update gui fs cache regularly
setInterval(puter.checkAndUpdateGUIFScache, 10000);
// initialize loop for updating caches for major directories; the GUI
// sets the token many times, but one loop is enough
if (this.env === 'gui' && !this.guiFSCacheTimer_) {
this.guiFSCacheTimer_ = setInterval(
() => this.checkAndUpdateGUIFScache(),
10000,
);
}
this._emitAuthStateChanged();
@@ -1907,6 +1911,14 @@ export class Puter {
}
};
/**
* The loop running `checkAndUpdateGUIFScache`, once a token is set.
*
* @internal
* @type {ReturnType<typeof setInterval> | null}
*/
guiFSCacheTimer_ = null;
/**
* Checks and updates the GUI FS cache for most-commonly used paths
*
@@ -2011,41 +2023,10 @@ if (puterParent) {
console.log('I have a parent, registering tools');
puterParent.on('message', async (event) => {
console.log('Got tool req ', event);
if (event.$ === 'requestTools') {
console.log('Responding with tools');
puterParent.postMessage({
$: 'providedTools',
tools: JSON.parse(JSON.stringify(puter.tools)),
});
}
if (event.$ === 'executeTool') {
console.log('xecuting tools');
/**
* Puter tools format
*
* @type {[
* {
* exec: Function;
* function: {
* description: string;
* name: string;
* parameters: { properties: any; required: string[] };
* type: string;
* };
* },
* ]}
*/
const [tool] = puter.tools.filter(
(e) => e.function.name === event.toolName,
);
const response = await tool.exec(event.parameters);
puterParent.postMessage({
$: 'toolResponse',
response,
tag: event.tag,
});
try {
await handleToolRequest(puter.tools, puterParent, event);
} catch (e) {
console.error(e);
}
});
puterParent.postMessage({ $: 'ready' });
@@ -2055,6 +2036,7 @@ globalThis.addEventListener &&
globalThis.addEventListener('message', async (event) => {
// if the message is not from Puter, then ignore it
if (event.origin !== puter.defaultGUIOrigin) return;
if (!event.data || typeof event.data !== 'object') return;
if (event.data.msg && event.data.msg === 'requestOrigin') {
event.source.postMessage(
@@ -2084,9 +2066,14 @@ globalThis.addEventListener &&
// Call onAuth callback
if (puter.onAuth && typeof puter.onAuth === 'function') {
puter.getUser().then((user) => {
puter.onAuth(user);
});
// Not awaited: the waiting sign-ins below settle first.
(async () => {
try {
puter.onAuth(await puter.getUser());
} catch (e) {
console.error(e);
}
})();
}
puter.puterAuthState.isPromptOpen = false;
+68
View File
@@ -285,3 +285,71 @@ describe('web-mode stored token', () => {
expect(location.search).toBe('?auth_token=sites-own-param');
});
});
describe('GUI cache refresh', () => {
it('keeps one refresh loop however often the token is set', async () => {
const puter = await bootApp(`${LAUNCH}${APP_TOKEN}`);
puter.env = 'gui';
const refresh = vi
.spyOn(puter, 'checkAndUpdateGUIFScache')
.mockImplementation(() => {});
vi.useFakeTimers({ toFake: ['setInterval', 'clearInterval'] });
try {
puter.setAuthToken(APP_TOKEN);
puter.setAuthToken(APP_TOKEN);
puter.setAuthToken(APP_TOKEN);
vi.advanceTimersByTime(10_000);
expect(refresh).toHaveBeenCalledOnce();
} finally {
vi.useRealTimers();
}
});
});
describe('GUI message listener', () => {
/** Boots an app and returns the window listener it adds last. */
const bootWithListener = async () => {
const add = vi.spyOn(window, 'addEventListener');
try {
const puter = await bootApp(`${LAUNCH}${APP_TOKEN}`);
const [, handler] = add.mock.calls
.filter(([type]) => type === 'message')
.at(-1);
return { puter, handler };
} finally {
add.mockRestore();
}
};
it('ignores a GUI message with no data', async () => {
const { puter, handler } = await bootWithListener();
await expect(
handler({ origin: puter.defaultGUIOrigin, data: null }),
).resolves.toBeUndefined();
});
it('reports a failed user lookup after sign-in', async () => {
const { puter, handler } = await bootWithListener();
const error = vi.spyOn(console, 'error').mockImplementation(() => {});
const unhandled = vi.fn();
process.on('unhandledRejection', unhandled);
try {
puter.onAuth = vi.fn();
puter.getUser = vi.fn(async () => {
throw new Error('offline');
});
await handler({
origin: puter.defaultGUIOrigin,
source: globalThis.parent,
data: { msg: 'puter.token', token: APP_TOKEN },
});
await new Promise((resolve) => setTimeout(resolve, 0));
expect(unhandled).not.toHaveBeenCalled();
expect(puter.onAuth).not.toHaveBeenCalled();
expect(error).toHaveBeenCalled();
} finally {
process.off('unhandledRejection', unhandled);
error.mockRestore();
}
});
});
+10 -14
View File
@@ -15,20 +15,15 @@ export default class EventListener {
// Array of all supported event names.
#eventNames;
// Map of eventName -> array of listeners
#eventListeners;
/** @type {Map<string, Function[]>} */
#eventListeners = new Map();
/** @param {(keyof EventMap & string)[] | string[]} eventNames */
constructor (eventNames) {
this.#eventNames = eventNames;
this.#eventListeners = (() => {
const map = new Map();
for ( let eventName of this.#eventNames ) {
map[eventName] = [];
}
return map;
})();
for ( const eventName of this.#eventNames ) {
this.#eventListeners.set(eventName, []);
}
}
/**
@@ -44,9 +39,10 @@ export default class EventListener {
console.error(`Event name '${eventName}' not supported`);
return;
}
this.#eventListeners[eventName].forEach((listener) => {
// A snapshot, so a handler that removes itself doesn't skip the next.
for ( const listener of [...this.#eventListeners.get(eventName)] ) {
listener(data);
});
}
}
/**
@@ -63,7 +59,7 @@ export default class EventListener {
console.error(`Event name '${eventName}' not supported`);
return;
}
this.#eventListeners[eventName].push(callback);
this.#eventListeners.get(eventName).push(callback);
return this;
}
@@ -80,7 +76,7 @@ export default class EventListener {
console.error(`Event name '${eventName}' not supported`);
return;
}
const listeners = this.#eventListeners[eventName];
const listeners = this.#eventListeners.get(eventName);
const index = listeners.indexOf(callback);
if ( index !== -1 ) {
listeners.splice(index, 1);
@@ -0,0 +1,56 @@
import { describe, expect, it, vi } from 'vitest';
import EventListener from './EventListener.js';
describe('EventListener', () => {
it('calls the next handler when one removes itself mid-emit', () => {
const emitter = new EventListener(['tick']);
const calls = [];
const once = () => {
calls.push('once');
emitter.off('tick', once);
};
emitter.on('tick', once);
emitter.on('tick', () => calls.push('second'));
emitter.emit('tick');
emitter.emit('tick');
expect(calls).toEqual(['once', 'second', 'second']);
});
it('runs a handler added mid-emit from the next emit on', () => {
const emitter = new EventListener(['tick']);
const late = vi.fn();
emitter.on('tick', () => emitter.on('tick', late));
emitter.emit('tick');
expect(late).not.toHaveBeenCalled();
emitter.emit('tick');
expect(late).toHaveBeenCalledOnce();
});
it('accepts event names that collide with Map members', () => {
const emitter = new EventListener(['size', 'get', 'constructor']);
const handler = vi.fn();
emitter.on('size', handler);
emitter.on('get', handler);
emitter.on('constructor', handler);
emitter.emit('size', 1);
emitter.emit('get', 2);
emitter.emit('constructor', 3);
expect(handler.mock.calls).toEqual([[1], [2], [3]]);
});
it('reports and ignores an unsupported event', () => {
const error = vi.spyOn(console, 'error').mockImplementation(() => {});
const emitter = new EventListener(['tick']);
expect(emitter.on('other', () => {})).toBeUndefined();
expect(emitter.off('other', () => {})).toBeUndefined();
emitter.emit('other');
expect(error).toHaveBeenCalledTimes(3);
error.mockRestore();
});
});
+12
View File
@@ -0,0 +1,12 @@
/**
* Escapes text for use inside HTML element content or a quoted attribute.
*
* @param {unknown} text
* @returns {string}
*/
export const escapeHtml = (text) =>
String(text)
.replace(/&/g, '&amp;')
.replace(/</g, '&lt;')
.replace(/>/g, '&gt;')
.replace(/"/g, '&quot;');
+155 -29
View File
@@ -158,12 +158,85 @@ async function resolveReauth(resp, { interactive = true, sentToken } = {}) {
return null;
}
// -- Idle timeout --
// A request that makes no progress for its limit is aborted and fails with
// `request_timeout`; any progress restarts the clock. Read-safe requests get
// the short limit only while progress is observable: before headers, or once
// body progress has been seen (the XHR shim reports none for a buffered body).
const READ_IDLE_TIMEOUT_MS = 60_000;
const IDLE_TIMEOUT_MS = 15 * 60_000;
const requestTimeoutError = () => ({
message: 'Request timed out.',
code: 'request_timeout',
});
/**
* Starts the idle clock for a request about to be sent. On expiry the XHR is
* flagged `_puterTimedOut` and aborted, so its `abort` listeners can tell a
* timeout from a cancellation.
*
* @param {XMLHttpRequest} xhr
* @param {{ timeout?: number }} spec - `timeout` replaces both limits; `0`
* turns the clock off.
* @param {boolean} readSafe - Eligible for the short limit.
* @param {boolean} watchUpload - Count upload progress too.
* @returns {() => void} Stops the clock.
*/
function watchIdle(xhr, spec, readSafe, watchUpload) {
if (spec.timeout === 0) return () => {};
let timer;
let stopped = false;
let bodyProgress = false;
const limit = () => {
if (spec.timeout !== undefined) return spec.timeout;
if (!readSafe) return IDLE_TIMEOUT_MS;
if (xhr.readyState < 2) return READ_IDLE_TIMEOUT_MS;
if (isNdjson(xhr.getResponseHeader('content-type'))) {
return IDLE_TIMEOUT_MS;
}
return bodyProgress ? READ_IDLE_TIMEOUT_MS : IDLE_TIMEOUT_MS;
};
const stop = () => {
stopped = true;
clearTimeout(timer);
};
const arm = () => {
if (stopped) return;
clearTimeout(timer);
timer = setTimeout(() => {
stop();
xhr._puterTimedOut = true;
xhr.abort();
}, limit());
};
const onBody = () => {
bodyProgress = true;
arm();
};
xhr.addEventListener('readystatechange', () => {
// DONE comes before load/error/abort; the XHR shim reaches it on a
// failed fetch, where its headers can't be read.
if (xhr.readyState === 4) return stop();
if (xhr.readyState === 3) bodyProgress = true;
arm();
});
xhr.addEventListener('progress', onBody);
if (watchUpload) xhr.upload?.addEventListener('progress', arm);
for (const type of ['load', 'error', 'abort', 'timeout']) {
xhr.addEventListener(type, stop);
}
arm();
return stop;
}
/**
* The one XHR builder both `initXhr` (utils.js) and `fetchUrl` wrap. Opens the
* request, applies headers/credentials/responseType, and stashes the whole
* `spec` on `xhr._puterReq` as the single replay representation — any attempt
* (reauth, transient) rebuilds it by calling `buildXhr(spec)` again, which
* re-reads the live token when `includePuterAuth`.
* re-reads the live token when `includePuterAuth`. Sending it starts the idle
* clock (see `watchIdle`).
*
* @param {Object} spec
* @param {string} spec.url - Full request URL.
@@ -177,9 +250,13 @@ async function resolveReauth(resp, { interactive = true, sentToken } = {}) {
* @param {boolean} [spec.withCredentials=true] Default is `true`
* @param {string} [spec.responseType=''] Default is `''`
* @param {Object} [spec.logId] - Pre-built apiCallLogger request id.
* @param {number} [spec.timeout] - Idle limit in ms, replacing the defaults;
* `0` turns it off.
* @param {{ readSafe?: boolean }} [opts] - `readSafe` selects the short idle
* limit.
* @returns {XMLHttpRequest}
*/
function buildXhr(spec) {
function buildXhr(spec, { readSafe = false } = {}) {
const {
url,
method = 'GET',
@@ -215,7 +292,16 @@ function buildXhr(spec) {
const origSend = xhr.send.bind(xhr);
xhr.send = function (body) {
spec.body = body;
return origSend(body);
// An upload listener makes the browser preflight a cross-origin
// request. One carrying a bearer is preflighted already; driver calls
// go out as simple requests and must stay that way.
const stopIdle = watchIdle(xhr, spec, readSafe, !!bearer);
try {
return origSend(body);
} catch (e) {
stopIdle();
throw e;
}
};
if (globalThis.puter?.apiCallLogger?.isEnabled()) {
@@ -314,7 +400,10 @@ async function parseResponse(xhr) {
}
const contentType = xhr.getResponseHeader('content-type');
if (contentType.startsWith('application/json')) {
// No declared type (a bodiless 204, a proxy that strips it): a success is
// opaque bytes, an error body is read like a JSON one.
const untypedError = !contentType && !(xhr.status >= 200 && xhr.status < 300);
if (contentType?.startsWith('application/json') || untypedError) {
const text = await xhr.response.text();
try {
return JSON.parse(text);
@@ -322,7 +411,7 @@ async function parseResponse(xhr) {
return text;
}
}
if (contentType.startsWith('application/octet-stream')) {
if (!contentType || contentType.startsWith('application/octet-stream')) {
return xhr.response;
}
return { success: true, result: xhr.response };
@@ -417,17 +506,17 @@ const sleep = (ms, signal) =>
return reject(
signal.reason ?? new DOMException('Aborted', 'AbortError'),
);
const t = setTimeout(resolve, ms);
signal?.addEventListener(
'abort',
() => {
clearTimeout(t);
reject(
signal.reason ?? new DOMException('Aborted', 'AbortError'),
);
},
{ once: true },
);
const onAbort = () => {
clearTimeout(t);
reject(signal.reason ?? new DOMException('Aborted', 'AbortError'));
};
// `once` only cleans up when abort fires; a caller reusing one signal
// across requests would otherwise collect a listener per retry.
const t = setTimeout(() => {
signal?.removeEventListener('abort', onAbort);
resolve();
}, ms);
signal?.addEventListener('abort', onAbort, { once: true });
});
const retryDelay = (attempt) => RETRY_DELAYS_MS[attempt - 1];
@@ -520,13 +609,16 @@ async function resolveVerificationGate(code, factors) {
/**
* Send one attempt. Resolves with a terminal outcome: { streamed: true, xhr,
* lineStream } — NDJSON, resolved at HEADERS_RECEIVED { xhr, status } —
* buffered response (any HTTP status) { networkError: true, xhr } — transport
* error Rejects only on abort. Per-line semantics (usage/email prompts,
* `toString`) belong to the caller's `shapeStream`.
* buffered response (any HTTP status) { networkError: true, timedOut?, xhr } —
* transport error or idle timeout. Rejects only on abort. Per-line semantics
* (usage/email prompts, `toString`) belong to the caller's `shapeStream`.
*
* @param {Object} spec
* @param {boolean} [readSafe=false] - Selects the short idle limit.
*/
function sendOnce(spec) {
function sendOnce(spec, readSafe = false) {
return new Promise((resolve, reject) => {
const xhr = buildXhr(spec);
const xhr = buildXhr(spec, { readSafe });
let streamed = false;
let responseComplete = false;
@@ -609,16 +701,17 @@ function sendOnce(spec) {
responseComplete = true;
signalStreamUpdate?.();
});
xhr.addEventListener('timeout', () => {
const error = { message: 'Network request timed out.', code: 'network_error' };
failStream(error);
resolve({ networkError: true, xhr });
});
const timedOut = () => {
failStream(requestTimeoutError());
resolve({ networkError: true, timedOut: true, xhr });
};
xhr.addEventListener('timeout', timedOut);
xhr.addEventListener('error', () => {
failStream({ message: 'Network request failed.', code: 'network_error' });
resolve({ networkError: true, xhr });
});
xhr.addEventListener('abort', () => {
if (xhr._puterTimedOut) return timedOut();
const error = spec.signal?.reason ??
new DOMException('Aborted', 'AbortError');
failStream(error);
@@ -654,6 +747,14 @@ function sendOnce(spec) {
async function classifyRetry(outcome, ctx) {
if (outcome.streamed) return null; // committed stream — never retried
// An idle timeout replays once, on the read-safe rule: a fresh connection
// is the cure for a dead one, and a stall that repeats isn't transient.
if (outcome.timedOut) {
if (ctx.done.has('timeout')) return null;
const decision = transientRetry(ctx);
if (decision) ctx.done.add('timeout');
return decision;
}
if (outcome.networkError) return transientRetry(ctx);
const { xhr, status } = outcome;
@@ -719,10 +820,11 @@ async function classifyRetry(outcome, ctx) {
* outcome, and retries on reauth / transient causes; otherwise hands the
* outcome to `shape`.
*
* @param {Object} spec - BuildXhr spec (+ optional buildBody, signal).
* @param {Object} spec - BuildXhr spec (+ optional buildBody, signal,
* timeout).
* @param {Object} opts
* @param {boolean} [opts.retrySafe=false] - Eligible for transient backoff
* retry. Default is `false`
* retry and the short idle limit. Default is `false`
* @param {boolean} [opts.retryGated=true] - Eligible for 429 backoff retry,
* regardless of method (the gate rejects before the handler runs). Default is
* `true`
@@ -746,7 +848,7 @@ async function sendWithRetry(
};
while (true) {
ctx.attempt++;
const outcome = await sendOnce(spec);
const outcome = await sendOnce(spec, retrySafe);
if (outcome.streamed)
return shapeStream(outcome.lineStream, outcome.xhr);
const decision = await classifyRetry(outcome, ctx);
@@ -833,6 +935,8 @@ function dedupe(key, factory, { windowMs = 2000 } = {}) {
* 401 may raise sign-in UI. Pass `false` for requests the user didn't ask for
* (boot telemetry, cache warmers): the stale token is dropped silently and
* the 401 surfaces to the caller instead. Default is `true`
* @param {number} [opts.timeout] - Idle limit in ms, replacing the defaults
* (see `watchIdle`); `0` turns it off.
* @param {Object} [opts.paginate] - Reserved for a later sprint step (ignored).
* @returns {Promise<PuterResponse>}
*/
@@ -850,6 +954,7 @@ function fetchUrl(url, opts = {}) {
retry,
dedupe: dedupeOpt,
interactiveReauth = true,
timeout,
} = opts;
const logId = logContext ?? {
@@ -869,6 +974,7 @@ function fetchUrl(url, opts = {}) {
signal,
logId,
interactiveReauth,
timeout,
};
// Read-safety: idempotent methods auto-retry; a POST read opts in with
@@ -889,6 +995,16 @@ function fetchUrl(url, opts = {}) {
return makeResponse(xhr, lineStream);
},
shape: async (outcome) => {
if (outcome.timedOut) {
if (loggingOn())
logRequest(logId, { error: requestTimeoutError() });
// A TypeError, as `fetch` rejects with, carrying the code.
const error = new TypeError(
`Network request to ${url} timed out`,
);
error.code = 'request_timeout';
throw error;
}
if (outcome.networkError) {
if (loggingOn())
logRequest(logId, {
@@ -1068,9 +1184,11 @@ function driverLineStream(lineStream, puter, upgradePrompt) {
* transform?: (result: unknown) => unknown;
* onError?: (error: unknown) => void;
* upgradePrompt?: UpgradePromptContext;
* timeout?: number;
* }} [opts]
* `readonly` marks the method retry-safe on transient failures (a
* rate/concurrency 429 replays either way — see GATE_REJECT_STATUS),
* `timeout` replaces the idle limits in ms (`0` turns it off),
* `transform` post-processes a successful result, `onError` is the legacy
* error callback the module APIs accept alongside the promise, and
* `upgradePrompt` is how the upgrade prompt names this method (defaulting to
@@ -1084,6 +1202,7 @@ async function driverCall(call, opts = {}) {
transform,
onError,
upgradePrompt,
timeout,
} = opts;
const puter = callInstance(call);
const promptContext = {
@@ -1118,6 +1237,7 @@ async function driverCall(call, opts = {}) {
responseType,
// Rebuilt per attempt, so a reauth replay carries the fresh token.
buildBody: () => callBody(call, puter),
timeout,
};
return await sendWithRetry(spec, {
@@ -1127,6 +1247,11 @@ async function driverCall(call, opts = {}) {
// Reauth and transient retries are already spent by the time the
// engine hands the outcome over, so this is terminal.
shape: async (outcome) => {
if (outcome.timedOut) {
const error = requestTimeoutError();
logCall(call, { error });
return fail(error);
}
if (outcome.networkError) {
logCall(call, { error: { message: 'Network error occurred' } });
return fail(outcome.xhr);
@@ -1225,6 +1350,7 @@ export {
fetchUrl,
isVerificationGateCode,
parseResponse,
requestTimeoutError,
resolveReauth,
resolveVerificationGate,
sendWithRetry,
+369 -6
View File
@@ -5,6 +5,7 @@ import {
driverCallEnvelope,
driverLineStream,
fetchUrl,
parseResponse,
sendWithRetry,
} from './networkUtils.js';
@@ -28,6 +29,7 @@ function installFakeXHR(program) {
this._respHeaders = {};
this.onreadystatechange = null;
this.onprogress = null;
this.upload = { addEventListener: vi.fn() };
instances.push(this);
}
open(method, url) {
@@ -54,18 +56,23 @@ function installFakeXHR(program) {
for (const [k, v] of Object.entries(headers))
this._respHeaders[k.toLowerCase()] = v;
}
_headersReceived() {
this.readyState = 2;
_setReadyState(state) {
if (this.readyState === state) return;
this.readyState = state;
this.onreadystatechange?.();
this.dispatchEvent(new Event('readystatechange'));
}
_headersReceived() {
this._setReadyState(2);
}
_progress(chunk) {
this.responseText += chunk;
this.readyState = 3;
this._setReadyState(3);
this.onprogress?.();
this.dispatchEvent(new Event('progress'));
}
_done() {
this.readyState = 4;
this.onreadystatechange?.();
this._setReadyState(4);
this.dispatchEvent(new Event('load'));
}
_networkError() {
@@ -684,6 +691,348 @@ describe('transient retry', () => {
});
});
describe('idle timeout', () => {
const READ_MS = 60_000;
const LONG_MS = 15 * 60_000;
const silent = () => {}; // never answers
const settledWith = async (promise) => {
try {
return { value: await promise };
} catch (error) {
return { error };
}
};
const PENDING = Symbol('pending');
const peek = (promise) => Promise.race([promise, Promise.resolve(PENDING)]);
beforeEach(() => {
globalThis.puter = {
authToken: 'tok',
APIOrigin: 'https://api.example',
env: 'nodejs',
};
vi.useFakeTimers();
});
afterEach(() => {
vi.useRealTimers();
});
it('times out a silent read after 60 s and replays it once', async () => {
const xhrs = installFakeXHR(silent);
const result = settledWith(fetchUrl('https://api.example/x'));
await vi.advanceTimersByTimeAsync(READ_MS - 1);
expect(await peek(result)).toBe(PENDING);
await vi.advanceTimersByTimeAsync(1);
await vi.advanceTimersByTimeAsync(250); // backoff before the replay
expect(xhrs.length).toBe(2);
await vi.advanceTimersByTimeAsync(READ_MS + 10_000);
const { error } = await result;
expect(error).toBeInstanceOf(TypeError);
expect(error.code).toBe('request_timeout');
expect(xhrs.length).toBe(2);
expect(vi.getTimerCount()).toBe(0);
});
it('gives a write 15 min and never replays it', async () => {
const xhrs = installFakeXHR(silent);
const result = settledWith(
fetchUrl('https://api.example/x', { method: 'POST', body: '{}' }),
);
await vi.advanceTimersByTimeAsync(LONG_MS - 1);
expect(await peek(result)).toBe(PENDING);
await vi.advanceTimersByTimeAsync(1);
const { error } = await result;
expect(error).toMatchObject({ code: 'request_timeout' });
await vi.advanceTimersByTimeAsync(LONG_MS);
expect(xhrs.length).toBe(1);
});
it('rejects a driver call with request_timeout and never replays it', async () => {
const xhrs = installFakeXHR(silent);
const result = settledWith(
driverCall({ iface: 'puter-chat-completion', method: 'complete', args: {} }),
);
await vi.advanceTimersByTimeAsync(LONG_MS - 1);
expect(await peek(result)).toBe(PENDING);
await vi.advanceTimersByTimeAsync(1);
expect((await result).error).toEqual({
message: 'Request timed out.',
code: 'request_timeout',
});
await vi.advanceTimersByTimeAsync(LONG_MS);
expect(xhrs.length).toBe(1);
});
it('replays a timed-out readonly driver call once', async () => {
const xhrs = installFakeXHR(silent);
const result = settledWith(
driverCall(
{ iface: 'puter-kvstore', method: 'get', args: { key: 'k' } },
{ readonly: true },
),
);
await vi.advanceTimersByTimeAsync(2 * READ_MS + 10_000);
expect((await result).error).toMatchObject({ code: 'request_timeout' });
expect(xhrs.length).toBe(2);
});
it('keeps a read alive while its body makes progress', async () => {
const xhrs = installFakeXHR(silent);
const result = settledWith(fetchUrl('https://api.example/x'));
const [xhr] = xhrs;
await vi.advanceTimersByTimeAsync(READ_MS - 10_000);
xhr._setHeaders(200, { 'content-type': 'application/json' });
xhr._headersReceived();
for (let i = 0; i < 5; i++) {
xhr._progress(' ');
await vi.advanceTimersByTimeAsync(READ_MS - 10_000);
}
xhr.responseText = '{"ok":true}';
xhr._done();
expect((await result).value.status).toBe(200);
expect(xhrs.length).toBe(1);
});
it('times out a read whose body stalls once progress was seen', async () => {
const xhrs = installFakeXHR((xhr) => {
xhr._setHeaders(200, { 'content-type': 'application/json' });
xhr._headersReceived();
xhr._progress('{');
});
const result = settledWith(fetchUrl('https://api.example/x'));
await vi.advanceTimersByTimeAsync(READ_MS);
expect(xhrs[0]._puterTimedOut).toBe(true);
await vi.advanceTimersByTimeAsync(250 + READ_MS);
expect((await result).error.code).toBe('request_timeout');
});
// The node/workerd shim reports no progress for a buffered body, so a
// read can't be held to 60 s once its headers are in.
it('moves a read to 15 min after headers until body progress is seen', async () => {
const xhrs = installFakeXHR((xhr) => {
xhr._setHeaders(200, { 'content-type': 'application/octet-stream' });
xhr._headersReceived();
});
const result = settledWith(fetchUrl('https://api.example/big'));
await vi.advanceTimersByTimeAsync(LONG_MS - 1);
expect(await peek(result)).toBe(PENDING);
expect(xhrs.length).toBe(1);
await vi.advanceTimersByTimeAsync(1);
expect(xhrs[0]._puterTimedOut).toBe(true);
// The replay finishes this time.
await vi.advanceTimersByTimeAsync(250);
xhrs[1].responseText = 'bytes';
xhrs[1]._done();
expect((await result).value.status).toBe(200);
});
it('gives a stream 15 min of silence before failing it', async () => {
installFakeXHR((xhr) => {
xhr._setHeaders(200, { 'content-type': 'application/x-ndjson' });
xhr._headersReceived();
xhr._progress('{"text":"a"}\n');
});
const response = await fetchUrl('https://api.example/stream');
const stream = response.stream();
expect((await stream.next()).value.text).toBe('a');
const read = settledWith(stream.next());
await vi.advanceTimersByTimeAsync(LONG_MS - 1);
expect(await peek(read)).toBe(PENDING);
await vi.advanceTimersByTimeAsync(1);
expect((await read).error).toEqual({
message: 'Request timed out.',
code: 'request_timeout',
});
});
it('lets upload progress reset the clock on a request with a bearer', async () => {
const xhrs = installFakeXHR(silent);
const result = settledWith(
fetchUrl('https://api.example/upload', {
method: 'POST',
includePuterAuth: true,
body: 'payload',
timeout: 1000,
}),
);
const [, onUpload] = xhrs[0].upload.addEventListener.mock.calls[0];
await vi.advanceTimersByTimeAsync(900);
onUpload();
await vi.advanceTimersByTimeAsync(900);
expect(await peek(result)).toBe(PENDING);
await vi.advanceTimersByTimeAsync(100);
expect((await result).error.code).toBe('request_timeout');
});
// A listener on `xhr.upload` would turn the driver call's simple CORS
// request into a preflighted one.
it('leaves upload progress alone on a driver call', async () => {
const xhrs = installFakeXHR(
respond({ body: { success: true, result: 1 } }),
);
await driverCall({ iface: 'puter-kvstore', method: 'get', args: {} });
expect(xhrs[0].upload.addEventListener).not.toHaveBeenCalled();
});
it('turns the clock off with timeout: 0', async () => {
installFakeXHR(silent);
const result = settledWith(
fetchUrl('https://api.example/x', { timeout: 0 }),
);
await vi.advanceTimersByTimeAsync(2 * LONG_MS);
expect(await peek(result)).toBe(PENDING);
expect(vi.getTimerCount()).toBe(0);
});
// The XHR shim marks a failed fetch DONE before dispatching `error`, and
// has no headers to read then.
it('settles a transport failure that reaches DONE without headers', async () => {
globalThis.puter.config = { autoRetry: false };
const xhrs = installFakeXHR((xhr) => {
xhr.getResponseHeader = () => {
throw new TypeError('no headers');
};
xhr._setReadyState(4);
xhr._networkError();
});
const result = settledWith(fetchUrl('https://api.example/x'));
expect((await result).error).toBeInstanceOf(TypeError);
expect(xhrs.length).toBe(1);
expect(vi.getTimerCount()).toBe(0);
});
it('clears the clock when a request settles', async () => {
installFakeXHR(respond({ body: { ok: true } }));
await fetchUrl('https://api.example/x');
expect(vi.getTimerCount()).toBe(0);
});
it('clears the clock when a request is cancelled', async () => {
installFakeXHR(silent);
const controller = new AbortController();
const result = settledWith(
fetchUrl('https://api.example/x', { signal: controller.signal }),
);
controller.abort();
expect((await result).error).toMatchObject({ name: 'AbortError' });
expect(vi.getTimerCount()).toBe(0);
});
});
describe('abort listeners on a reused signal', () => {
/** An AbortSignal that counts its live `abort` listeners. */
const trackedSignal = () => {
const signal = new AbortController().signal;
const live = new Set();
const add = signal.addEventListener.bind(signal);
const remove = signal.removeEventListener.bind(signal);
signal.addEventListener = (type, fn, opts) => {
if (type === 'abort') live.add(fn);
add(type, fn, opts);
};
signal.removeEventListener = (type, fn, opts) => {
if (type === 'abort') live.delete(fn);
remove(type, fn, opts);
};
return { signal, live };
};
beforeEach(() => {
globalThis.puter = {};
});
it('removes every listener once a retried request settles', async () => {
vi.useFakeTimers();
const { signal, live } = trackedSignal();
installFakeXHR(
sequence(
respond({ status: 503, body: {} }),
respond({ status: 503, body: {} }),
respond({ status: 200, body: {} }),
respond({ status: 503, body: {} }),
respond({ status: 200, body: {} }),
),
);
for (let i = 0; i < 2; i++) {
const p = fetchUrl('https://api.example/x', { signal });
await vi.advanceTimersByTimeAsync(60_000);
expect((await p).status).toBe(200);
}
expect(live.size).toBe(0);
vi.useRealTimers();
});
it('still aborts a retry wait', async () => {
vi.useFakeTimers();
const controller = new AbortController();
const xhrs = installFakeXHR(respond({ status: 503, body: {} }));
const p = fetchUrl('https://api.example/x', {
signal: controller.signal,
});
let error;
try {
await vi.advanceTimersByTimeAsync(100);
controller.abort();
await p;
} catch (e) {
error = e;
}
expect(error).toMatchObject({ name: 'AbortError' });
expect(xhrs.length).toBe(1);
vi.useRealTimers();
});
});
describe('parseResponse', () => {
const blobXhr = ({ status = 200, contentType = null, body = '' }) => ({
responseType: 'blob',
status,
response: new Blob([body]),
getResponseHeader: (name) =>
name.toLowerCase() === 'content-type' ? contentType : null,
});
it('returns the body of a success that declares no content type', async () => {
const xhr = blobXhr({ body: 'bytes' });
expect(await parseResponse(xhr)).toBe(xhr.response);
});
it('reads an error body that declares no content type', async () => {
expect(
await parseResponse(
blobXhr({ status: 502, body: '{"code":"bad_gateway"}' }),
),
).toEqual({ code: 'bad_gateway' });
expect(
await parseResponse(blobXhr({ status: 502, body: 'Bad Gateway' })),
).toBe('Bad Gateway');
});
it('keeps the typed blob branches', async () => {
const octet = blobXhr({ contentType: 'application/octet-stream' });
expect(await parseResponse(octet)).toBe(octet.response);
const png = blobXhr({ contentType: 'image/png' });
expect(await parseResponse(png)).toEqual({
success: true,
result: png.response,
});
});
});
describe('dedupe', () => {
it('coalesces concurrent identical requests into one call', async () => {
let calls = 0;
@@ -757,6 +1106,20 @@ describe('driverCall', () => {
});
});
it('resolves a blob response that declares no content type', async () => {
installFakeXHR((xhr) => {
xhr._setHeaders(200);
xhr._headersReceived();
xhr.response = new Blob(['image']);
xhr._done();
});
const result = await driverCall(
{ iface: 'puter-image-generation', method: 'generate', args: {} },
{ responseType: 'blob' },
);
expect(result).toBeInstanceOf(Blob);
});
it('resolves the whole response when the driver returns no result field', async () => {
installFakeXHR(respond({ body: { success: true, models: [] } }));
expect(await driverCall(call)).toEqual({ success: true, models: [] });
@@ -1148,7 +1511,7 @@ describe('NDJSON stream termination', () => {
const { xhr, stream } = await start();
const read = stream.next();
xhr.dispatchEvent(new Event('timeout'));
await expect(read).rejects.toMatchObject({ code: 'network_error' });
await expect(read).rejects.toMatchObject({ code: 'request_timeout' });
});
});
+56
View File
@@ -0,0 +1,56 @@
/** @typedef {import('./types.js').ToolSchema} ToolSchema */
/**
* Answers one tool request from the parent app. An `executeTool` request always
* gets a `toolResponse` back, with `error` set when the tool is unknown or
* fails, so the parent's wait settles.
*
* @param {ToolSchema[]} tools
* @param {{ postMessage: (message: unknown) => void }} parent
* @param {unknown} event
* @returns {Promise<void>}
*/
export async function handleToolRequest(tools, parent, event) {
if (!event || typeof event !== 'object') return;
const request = /** @type {Record<string, unknown>} */ (event);
if (request.$ === 'requestTools') {
console.log('Responding with tools');
parent.postMessage({
$: 'providedTools',
tools: JSON.parse(JSON.stringify(tools)),
});
return;
}
if (request.$ !== 'executeTool') return;
console.log('Executing tools');
const fail = (message, code) => {
parent.postMessage({
$: 'toolResponse',
tag: request.tag,
error: { message, code },
});
};
const tool = tools.find((t) => t?.function?.name === request.toolName);
if (!tool) {
fail(`Unknown tool: ${request.toolName}`, 'tool_not_found');
return;
}
let response;
try {
response = await tool.exec(
/** @type {Record<string, unknown>} */ (request.parameters),
);
} catch (e) {
fail(e?.message ?? String(e), 'tool_failed');
return;
}
try {
parent.postMessage({ $: 'toolResponse', response, tag: request.tag });
} catch (e) {
// e.g. a result the structured clone can't copy
fail(e?.message ?? String(e), 'tool_failed');
}
}
+92
View File
@@ -0,0 +1,92 @@
import { beforeEach, describe, expect, it, vi } from 'vitest';
import { handleToolRequest } from './toolBridge.js';
const tool = (name, exec) => ({
function: { name, description: name, parameters: {} },
exec,
});
let parent;
beforeEach(() => {
parent = { postMessage: vi.fn() };
vi.spyOn(console, 'log').mockImplementation(() => {});
});
describe('handleToolRequest', () => {
it('lists the tools', async () => {
const tools = [tool('add', () => 3)];
await handleToolRequest(tools, parent, { $: 'requestTools' });
expect(parent.postMessage).toHaveBeenCalledWith({
$: 'providedTools',
tools: [{ function: tools[0].function }],
});
});
it('answers with the tool result', async () => {
const exec = vi.fn(async ({ a, b }) => a + b);
await handleToolRequest([tool('add', exec)], parent, {
$: 'executeTool',
toolName: 'add',
parameters: { a: 1, b: 2 },
tag: 't1',
});
expect(parent.postMessage).toHaveBeenCalledWith({
$: 'toolResponse',
response: 3,
tag: 't1',
});
});
it('answers an unknown tool with an error', async () => {
await handleToolRequest([tool('add', () => 3)], parent, {
$: 'executeTool',
toolName: 'missing',
tag: 't2',
});
expect(parent.postMessage).toHaveBeenCalledWith({
$: 'toolResponse',
tag: 't2',
error: { message: 'Unknown tool: missing', code: 'tool_not_found' },
});
});
it('answers a failing tool with an error', async () => {
const exec = () => {
throw new Error('boom');
};
await handleToolRequest([tool('add', exec)], parent, {
$: 'executeTool',
toolName: 'add',
tag: 't3',
});
expect(parent.postMessage).toHaveBeenCalledWith({
$: 'toolResponse',
tag: 't3',
error: { message: 'boom', code: 'tool_failed' },
});
});
it('answers with an error when the result cannot be posted', async () => {
parent.postMessage.mockImplementationOnce(() => {
throw new Error('could not be cloned');
});
await handleToolRequest([tool('add', () => () => {})], parent, {
$: 'executeTool',
toolName: 'add',
tag: 't4',
});
expect(parent.postMessage).toHaveBeenLastCalledWith({
$: 'toolResponse',
tag: 't4',
error: { message: 'could not be cloned', code: 'tool_failed' },
});
});
it.each([null, undefined, 'text', { $: 'other' }])(
'ignores %s',
async (event) => {
await handleToolRequest([], parent, event);
expect(parent.postMessage).not.toHaveBeenCalled();
},
);
});
+17 -2
View File
@@ -1,7 +1,7 @@
import { FileReaderPoly } from './polyfills/fileReaderPoly.js';
import {
buildXhr, driverCall, isVerificationGateCode, parseResponse, resolveReauth,
resolveVerificationGate,
buildXhr, driverCall, isVerificationGateCode, parseResponse, requestTimeoutError,
resolveReauth, resolveVerificationGate,
} from './networkUtils.js';
/**
@@ -213,6 +213,21 @@ function setupXhrEventHandlers (xhr, success_cb, error_cb, resolve_func, reject_
}
return handle_error(error_cb, reject_func, this);
});
// The idle timeout ends a request by aborting it (see `watchIdle`).
xhr.addEventListener('abort', function () {
if ( ! xhr._puterTimedOut ) return;
const error = requestTimeoutError();
if ( globalThis.puter?.apiCallLogger?.isEnabled() && xhr._puterRequestId ) {
globalThis.puter.apiCallLogger.logRequest({
service: xhr._puterRequestId.service,
operation: xhr._puterRequestId.operation,
params: xhr._puterRequestId.params,
error,
});
}
return handle_error(error_cb, reject_func, error);
});
}
/**
+68 -1
View File
@@ -1,5 +1,5 @@
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest';
import { makeDriverMethod } from './utils.js';
import { initXhr, makeDriverMethod, setupXhrEventHandlers } from './utils.js';
/**
* Pins the callback contract of `makeDriverMethod`: driver methods are
@@ -114,3 +114,70 @@ describe('makeDriverMethod legacy callbacks', () => {
});
});
});
describe('setupXhrEventHandlers', () => {
it('settles a blob read whose response declares no content type', async () => {
const listeners = {};
const xhr = {
responseType: 'blob',
status: 200,
response: new Blob(['bytes']),
getResponseHeader: () => null,
addEventListener: (type, fn) => { listeners[type] = fn; },
};
const result = new Promise((resolve, reject) => {
setupXhrEventHandlers(xhr, undefined, undefined, resolve, reject);
});
listeners.load.call(xhr, { target: xhr });
const pending = new Promise(resolve => setTimeout(() => resolve('pending'), 100));
expect(await Promise.race([result, pending])).toBe(xhr.response);
});
});
describe('legacy request idle timeout', () => {
// An XHR that is sent but never answered.
class SilentXHR extends EventTarget {
readyState = 0;
upload = { addEventListener () {} };
open () { this.readyState = 1; }
setRequestHeader () {}
getResponseHeader () { return null; }
send () {}
abort () { this.dispatchEvent(new Event('abort')); }
}
beforeEach(() => {
vi.useFakeTimers();
globalThis.XMLHttpRequest = SilentXHR;
});
afterEach(() => {
vi.useRealTimers();
});
it.each(['post', 'get'])('rejects a silent %s with request_timeout after 15 min', async (method) => {
const onError = vi.fn();
let outcome;
const xhr = initXhr('/stat', 'https://api.test', 'tok', method);
const request = new Promise((resolve, reject) => {
setupXhrEventHandlers(xhr, undefined, onError, resolve, reject);
});
(async () => {
try {
outcome = { value: await request };
} catch (error) {
outcome = { error };
}
})();
xhr.send('{}');
await vi.advanceTimersByTimeAsync(15 * 60_000 - 1);
expect(outcome).toBeUndefined();
await vi.advanceTimersByTimeAsync(1);
const error = { message: 'Request timed out.', code: 'request_timeout' };
expect(outcome).toEqual({ error });
expect(onError).toHaveBeenCalledOnce();
expect(onError).toHaveBeenCalledWith(error);
expect(vi.getTimerCount()).toBe(0);
});
});
+6 -3
View File
@@ -54,15 +54,17 @@ export class CallbackManager {
/**
* Registers `callback` and binds it to `source`, the only window whose
* messages may invoke it later. A callback registered without a source
* can never be invoked from outside this document.
* can never be invoked from outside this document. A `once` callback is
* dropped after its first call; anything else lives as long as the page.
*
* @param {Function} callback
* @param {Window} [source]
* @param {{ once?: boolean }} [options]
* @returns {string}
*/
register_callback (callback, source) {
register_callback (callback, source, { once = false } = {}) {
const id = randomCallbackId();
this.callbacks.set(id, { callback, source });
this.callbacks.set(id, { callback, source, once });
return id;
}
@@ -80,6 +82,7 @@ export class CallbackManager {
// Only the window the callback was dehydrated for may invoke it,
// otherwise a sibling frame could drive another app's callbacks.
if ( ! entry || event.source !== entry.source ) return;
if ( entry.once ) this.callbacks.delete(data.id);
entry.callback(...(Array.isArray(data.args) ? data.args : []));
});
}
+54
View File
@@ -126,3 +126,57 @@ describe('xdrpc callback source binding', () => {
}
});
});
describe('xdrpc callback lifetime', () => {
const $SCOPE = '9a9c83a4-7897-43a0-93b9-53217b84fde6';
const setup = () => {
const handlers = [];
const manager = new CallbackManager();
manager.attach_to_source({
addEventListener: (type, handler) =>
type === 'message' && handlers.push(handler),
});
const deliver = event => handlers.forEach(handler => handler(event));
return { manager, deliver };
};
it('drops a one-shot callback after it fires', () => {
const gui = {};
const { manager, deliver } = setup();
const calls = [];
const id = manager.register_callback(v => calls.push(v), gui, { once: true });
deliver({ source: gui, data: { $SCOPE, id, args: ['first'] } });
deliver({ source: gui, data: { $SCOPE, id, args: ['second'] } });
expect(calls).toEqual(['first']);
expect(manager.callbacks.size).toBe(0);
});
it('keeps a one-shot callback a forged message failed to invoke', () => {
const gui = {};
const { manager, deliver } = setup();
const calls = [];
const id = manager.register_callback(v => calls.push(v), gui, { once: true });
deliver({ source: { sibling: true }, data: { $SCOPE, id, args: ['forged'] } });
deliver({ source: gui, data: { $SCOPE, id, args: ['real'] } });
expect(calls).toEqual(['real']);
});
it('keeps a dehydrated callback for repeat calls', () => {
const gui = {};
const { manager, deliver } = setup();
const calls = [];
const { id } = new Dehydrator({ callbackManager: manager, source: gui })
.dehydrate(() => calls.push('clicked'));
deliver({ source: gui, data: { $SCOPE, id, args: [] } });
deliver({ source: gui, data: { $SCOPE, id, args: [] } });
expect(calls).toEqual(['clicked', 'clicked']);
expect(manager.callbacks.size).toBe(1);
});
});
+2 -1
View File
@@ -25,12 +25,13 @@ export class Debug {
this.puter.logger.on(category);
}
globalThis.addEventListener('message', async e => {
globalThis.addEventListener('message', e => {
// Ensure message is from parent window
if ( e.source !== globalThis.parent ) return;
// (parent window is allowed to be anything)
// Check if it's a debug message
if ( ! e.data || typeof e.data !== 'object' ) return;
if ( ! e.data.$ ) return;
if ( e.data.$ !== 'puterjs-debug' ) return;
+55
View File
@@ -0,0 +1,55 @@
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest';
import { Debug } from './Debug.js';
let handler;
let parent;
beforeEach(() => {
parent = {};
vi.stubGlobal('location', { href: 'https://app.test/?enabled_logs=net' });
vi.stubGlobal('parent', parent);
vi.stubGlobal('addEventListener', (type, fn) => {
if (type === 'message') handler = fn;
});
vi.spyOn(console, 'log').mockImplementation(() => {});
});
afterEach(() => {
vi.unstubAllGlobals();
vi.restoreAllMocks();
});
const makeDebug = () => {
const on = vi.fn();
new Debug({ logger: { on } });
return on;
};
/** Runs the listener the way the browser does: nothing awaits it. */
const dispatch = (event) => Promise.resolve().then(() => handler(event));
describe('Debug', () => {
it('turns on the categories named in the URL', () => {
expect(makeDebug()).toHaveBeenCalledWith('net');
});
it.each([null, undefined, 'text'])(
'ignores a parent message carrying %s',
async (data) => {
const on = makeDebug();
await expect(
dispatch({ source: parent, data }),
).resolves.toBeUndefined();
expect(on).toHaveBeenCalledTimes(1);
},
);
it('turns a category on when the parent asks', async () => {
const on = makeDebug();
await dispatch({
source: parent,
data: { $: 'puterjs-debug', cmd: 'log.on', category: 'fs' },
});
expect(on).toHaveBeenCalledWith('fs');
});
});
@@ -1,3 +1,5 @@
import { escapeHtml } from '../lib/html.js';
class EmailConfirmationDialog extends (globalThis.HTMLElement || Object) {
constructor (message) {
super();
@@ -198,7 +200,7 @@ class EmailConfirmationDialog extends (globalThis.HTMLElement || Object) {
</svg>
</div>
<h2>Confirm Your Email</h2>
<p class="message">${this.message}</p>
<p class="message">${escapeHtml(this.message)}</p>
<div class="buttons">
<button class="button button-primary" id="confirm-email-btn">Go to Puter.com</button>
<button class="button button-cancel" id="close-btn">Close</button>
@@ -0,0 +1,38 @@
import { afterAll, describe, expect, it, vi } from 'vitest';
// Just enough of an element for the constructor to render into.
vi.stubGlobal(
'HTMLElement',
class {
attachShadow() {
this.shadowRoot = { innerHTML: '' };
return this.shadowRoot;
}
},
);
const { default: EmailConfirmationDialog } =
await import('./EmailConfirmationDialog.js');
afterAll(() => {
vi.unstubAllGlobals();
});
describe('EmailConfirmationDialog', () => {
it('renders the message as text', () => {
const dialog = new EmailConfirmationDialog(
'<img src=x onerror="alert(1)"> & more',
);
const html = dialog.shadowRoot.innerHTML;
expect(html).toContain(
'&lt;img src=x onerror=&quot;alert(1)&quot;&gt; &amp; more',
);
expect(html).not.toContain('<img');
});
it('falls back to the default message', () => {
const dialog = new EmailConfirmationDialog();
expect(dialog.shadowRoot.innerHTML).toContain(
'Please confirm your email address to use this service.',
);
});
});
@@ -59,21 +59,27 @@ afterEach(() => {
describe('cache update timer', () => {
it('keeps a single interval across repeated auth-state changes', () => {
const fs = makeModule();
// `vi.getTimerCount()` is global and constructing the module schedules
// a timer of its own, so the cache interval is counted as a delta from
// construction rather than as an absolute.
const baseline = vi.getTimerCount();
// Each auth change also sends the cache timestamp request, whose idle
// clock is a timer too, so intervals are counted on their own.
const started = vi.spyOn(globalThis, 'setInterval');
const cleared = vi.spyOn(globalThis, 'clearInterval');
const running = () => started.mock.calls.length - cleared.mock.calls.length;
fs.onAuthStateChanged();
const timer = fs.cacheUpdateTimer;
fs.onAuthStateChanged();
fs.onAuthStateChanged();
try {
fs.onAuthStateChanged();
const timer = fs.cacheUpdateTimer;
fs.onAuthStateChanged();
fs.onAuthStateChanged();
expect(vi.getTimerCount()).toBe(baseline + 1);
expect(fs.cacheUpdateTimer).not.toBe(timer);
expect(running()).toBe(1);
expect(fs.cacheUpdateTimer).not.toBe(timer);
fs.stopCacheUpdateTimer();
expect(vi.getTimerCount()).toBe(baseline);
fs.stopCacheUpdateTimer();
expect(running()).toBe(0);
} finally {
started.mockRestore();
cleared.mockRestore();
}
});
it('refreshes the cache timestamp while running', () => {
+29 -17
View File
@@ -349,7 +349,7 @@ export class AppConnection extends EventListener {
// TODO: Set this.#puterOrigin to the puter origin
(globalThis.document) && window.addEventListener('message', event => {
const onMessage = event => {
// Relayed by the host environment; a window that guessed an
// appInstanceID must not be able to forge one directly.
if ( event.source !== this.messageTarget ) return;
@@ -376,12 +376,15 @@ export class AppConnection extends EventListener {
}
this.#isOpen = false;
// Nothing arrives for a closed app; don't hold the window.
window.removeEventListener('message', onMessage);
this.emit('close', {
appInstanceID: this.targetAppInstanceID,
statusCode: event.data.statusCode,
});
}
});
};
(globalThis.document) && window.addEventListener('message', onMessage);
}
/**
@@ -597,6 +600,14 @@ export class UIModule extends EventListener {
};
}
// Runs and drops the callback waiting on reply `msgId`, if any.
#settleCallback (msgId, value) {
const callback = msgId ? this.#callbackFunctions[msgId] : undefined;
if ( ! callback ) return;
delete this.#callbackFunctions[msgId];
callback(value);
}
#postMessageAsync (name, args = {}) {
return new Promise(resolve => {
this.#postMessageWithCallback(name, resolve, args);
@@ -627,7 +638,7 @@ export class UIModule extends EventListener {
done_setting_resolve();
});
});
const callback_id = this.util.rpc.registerCallback(resolve, this.messageTarget);
const callback_id = this.util.rpc.registerCallback(resolve, this.messageTarget, { once: true });
this.messageTarget?.postMessage({
$: 'puter-ipc',
v: 2,
@@ -697,12 +708,12 @@ export class UIModule extends EventListener {
// Bind the message event listener to the window
let lastDraggedOverElement = null;
(globalThis.document) && window.addEventListener('message', async (e) => {
const handleMessage = async (e) => {
if ( ! this.#isTrustedMessageSource(e) ) return;
if ( ! e.data ) return;
// `error`
if ( e.data.error ) {
throw e.data.error;
console.error(e.data.error);
}
// `focus` event
else if ( e.data.msg && e.data.msg === 'focus' ) {
@@ -828,28 +839,20 @@ export class UIModule extends EventListener {
// getAppDataSucceeded
else if ( e.data.msg === 'getAppDataSucceeded' ) {
let appDataItem = new FSItem(e.data.item);
if ( e.data.original_msg_id && this.#callbackFunctions[e.data.original_msg_id] ) {
this.#callbackFunctions[e.data.original_msg_id](appDataItem);
}
this.#settleCallback(e.data.original_msg_id, appDataItem);
}
// instancesOpenSucceeded
else if ( e.data.msg === 'instancesOpenSucceeded' ) {
if ( e.data.original_msg_id && this.#callbackFunctions[e.data.original_msg_id] ) {
this.#callbackFunctions[e.data.original_msg_id](e.data.instancesOpen);
}
this.#settleCallback(e.data.original_msg_id, e.data.instancesOpen);
}
// readAppDataFileSucceeded
else if ( e.data.msg === 'readAppDataFileSucceeded' ) {
let appDataItem = new FSItem(e.data.item);
if ( e.data.original_msg_id && this.#callbackFunctions[e.data.original_msg_id] ) {
this.#callbackFunctions[e.data.original_msg_id](appDataItem);
}
this.#settleCallback(e.data.original_msg_id, appDataItem);
}
// readAppDataFileFailed
else if ( e.data.msg === 'readAppDataFileFailed' ) {
if ( e.data.original_msg_id && this.#callbackFunctions[e.data.original_msg_id] ) {
this.#callbackFunctions[e.data.original_msg_id](null);
}
this.#settleCallback(e.data.original_msg_id, null);
}
// Determine if this is a response to a previous message and if so, is there
// a callback function for this message? if answer is yes to both then execute the callback
@@ -984,6 +987,15 @@ export class UIModule extends EventListener {
conn, accept, reject,
});
}
};
// Nothing awaits a message listener, so a throw would only surface as
// an unhandled rejection.
(globalThis.document) && window.addEventListener('message', async (e) => {
try {
await handleMessage(e);
} catch (err) {
console.error(err);
}
});
// We need to send the mouse position to the host environment
@@ -0,0 +1,202 @@
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest';
import { $SCOPE } from '../lib/xdrpc.js';
import { UtilRPC } from './Util.js';
const { AppConnection, UIModule } = await import('./UI.js');
let messageHandlers;
let parentPostMessage;
let parentWindow;
const listen = (name, handler) => {
if (name === 'message') messageHandlers.push(handler);
};
/** Deliver one message event to every listener registered so far. */
const deliver = async (event) => {
for (const handler of [...messageHandlers]) await handler(event);
};
/** A logger whose every chain ends quietly. */
const quietLogger = () => {
const logger = { info: () => {}, fields: () => logger };
return logger;
};
const makePuter = () => ({
env: 'app',
authToken: 'token-1',
APIOrigin: 'https://api.test',
appID: 'app-1',
logger: quietLogger(),
util: { rpc: new UtilRPC() },
});
const makeUI = () => new UIModule(makePuter(), { appInstanceID: 'instance-1' });
beforeEach(() => {
messageHandlers = [];
parentPostMessage = vi.fn();
parentWindow = { postMessage: parentPostMessage };
globalThis.window = {
parent: parentWindow,
focus: vi.fn(),
addEventListener: listen,
removeEventListener: (name, handler) => {
if (name !== 'message') return;
const i = messageHandlers.indexOf(handler);
if (i !== -1) messageHandlers.splice(i, 1);
},
};
globalThis.addEventListener = listen;
globalThis.document = { addEventListener: vi.fn() };
globalThis.puter = { defaultGUIOrigin: 'https://puter.com' };
});
afterEach(() => {
delete globalThis.window;
delete globalThis.addEventListener;
delete globalThis.document;
delete globalThis.puter;
vi.restoreAllMocks();
vi.unstubAllGlobals();
});
describe('AppConnection', () => {
const connect = (puter = makePuter()) =>
new AppConnection(puter, {
target: 'child-1',
usesSDK: true,
messageTarget: parentWindow,
appInstanceID: 'instance-1',
});
const fromChild = (data) => ({
source: parentWindow,
data: {
appInstanceID: 'child-1',
targetAppInstanceID: 'instance-1',
...data,
},
});
it('stops listening once the target app closes', async () => {
const puter = makePuter();
const before = messageHandlers.length;
const conn = connect(puter);
expect(messageHandlers.length).toBe(before + 1);
const onClose = vi.fn();
const onMessage = vi.fn();
conn.on('close', onClose);
conn.on('message', onMessage);
await deliver(fromChild({ msg: 'messageToApp', contents: 'hi' }));
await deliver(fromChild({ msg: 'appClosed', statusCode: 0 }));
await deliver(fromChild({ msg: 'appClosed', statusCode: 0 }));
await deliver(fromChild({ msg: 'messageToApp', contents: 'late' }));
expect(onMessage.mock.calls).toEqual([['hi']]);
expect(onClose).toHaveBeenCalledOnce();
expect(onClose).toHaveBeenCalledWith({
appInstanceID: 'child-1',
statusCode: 0,
});
expect(messageHandlers.length).toBe(before);
});
it('keeps listening when another app closes', async () => {
const conn = connect();
const count = messageHandlers.length;
const onClose = vi.fn();
conn.on('close', onClose);
await deliver(
fromChild({ msg: 'appClosed', appInstanceID: 'child-2' }),
);
expect(onClose).not.toHaveBeenCalled();
expect(messageHandlers.length).toBe(count);
});
});
describe('host replies (env: app)', () => {
it('frees an IPC reply callback once the reply arrives', async () => {
const ui = makeUI();
const { callbackManager } = ui.util.rpc;
const before = callbackManager.callbacks.size;
const result = ui.exitPictureInPicture();
// The stub registers its callback after a tick.
await vi.waitFor(() => expect(parentPostMessage).toHaveBeenCalled());
const { uuid } = parentPostMessage.mock.calls.at(-1)[0];
await deliver({
source: parentWindow,
data: { $SCOPE, id: uuid, args: [{ wasOpen: true }] },
});
await expect(result).resolves.toBe(true);
expect(callbackManager.callbacks.size).toBe(before);
});
// A settled reply's id must not keep catching later messages ahead of the
// branches below the reply dispatch.
it.each([
['instancesOpenSucceeded', { instancesOpen: 2 }],
['getAppDataSucceeded', { item: { uid: 'u' } }],
['readAppDataFileSucceeded', { item: { uid: 'u' } }],
['readAppDataFileFailed', {}],
])('settles a %s reply once', async (msg, fields) => {
const ui = makeUI();
ui.instancesOpen();
const { uuid } = parentPostMessage.mock.calls.at(-1)[0];
const onLocale = vi.fn();
ui.on('localeChanged', onLocale);
await deliver({
source: parentWindow,
data: { msg, original_msg_id: uuid, ...fields },
});
await deliver({
source: parentWindow,
data: {
msg: 'broadcast',
name: 'localeChanged',
data: { language: 'fr' },
original_msg_id: uuid,
},
});
expect(onLocale).toHaveBeenCalledWith({ language: 'fr' });
});
it('reports an error reply instead of rejecting the listener', async () => {
const error = vi.spyOn(console, 'error').mockImplementation(() => {});
makeUI();
const failure = { code: 'method_removed', message: 'gone' };
await expect(
deliver({
source: parentWindow,
data: { msg: 'error', original_msg_id: 1, error: failure },
}),
).resolves.toBeUndefined();
expect(error).toHaveBeenCalledWith(failure);
});
it('survives a malformed message', async () => {
const error = vi.spyOn(console, 'error').mockImplementation(() => {});
vi.stubGlobal('location', { search: '' });
const ui = makeUI();
ui.onItemsOpened(() => {});
await expect(
deliver({
source: parentWindow,
data: { msg: 'itemsOpened', msg_id: 1 },
}),
).resolves.toBeUndefined();
expect(error).toHaveBeenCalled();
});
});
+2 -6
View File
@@ -1,10 +1,6 @@
/** @typedef {{ title?: string, method?: string }} UsageLimitDialogOptions */
import { escapeHtml } from '../lib/html.js';
const escapeHtml = (text) => String(text)
.replace(/&/g, '&amp;')
.replace(/</g, '&lt;')
.replace(/>/g, '&gt;')
.replace(/"/g, '&quot;');
/** @typedef {{ title?: string, method?: string }} UsageLimitDialogOptions */
class UsageLimitDialog extends (globalThis.HTMLElement || Object) {
/**
+5 -3
View File
@@ -52,14 +52,16 @@ export class UtilRPC {
}
/**
* Registers a function under a callback id `source` can invoke.
* Registers a function under a callback id `source` can invoke. Pass
* `once` for a reply, so the id is freed after its first call.
*
* @param {(value: unknown) => void} resolve
* @param {Window} [source]
* @param {{ once?: boolean }} [options]
* @returns {string}
*/
registerCallback (resolve, source) {
return this.callbackManager.register_callback(resolve, source);
registerCallback (resolve, source, options) {
return this.callbackManager.register_callback(resolve, source, options);
}
/**