Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
30 changes: 25 additions & 5 deletions src/client.ts
Original file line number Diff line number Diff line change
Expand Up @@ -100,8 +100,28 @@ function parseRetryAfterMs(res: Response): number | undefined {
return Number.isFinite(at) ? Math.max(0, at - Date.now()) : undefined;
}

function sleep(ms: number): Promise<void> {
return new Promise((resolve) => setTimeout(resolve, ms));
function abortError(signal?: AbortSignal | null): Error {
const reason = signal?.reason;
return reason instanceof Error
? reason
: new DOMException('The operation was aborted', 'AbortError');
}

function sleep(ms: number, signal?: AbortSignal | null): Promise<void> {
if (signal?.aborted) return Promise.reject(abortError(signal));
return new Promise((resolve, reject) => {
let timer: ReturnType<typeof setTimeout> | undefined;
const onAbort = () => {
if (timer !== undefined) clearTimeout(timer);
signal?.removeEventListener('abort', onAbort);
reject(abortError(signal));
};
timer = setTimeout(() => {
signal?.removeEventListener('abort', onAbort);
resolve();
}, ms);
signal?.addEventListener('abort', onAbort, { once: true });
});
}

export async function request<T>(
Expand Down Expand Up @@ -138,7 +158,7 @@ export async function request<T>(
if (isRetryableStatus(res.status) && attempt < retries) {
await res.text().catch(() => ''); // drain body so the socket can be reused
clearTimeout(timer);
await sleep(backoffDelayMs(attempt, retryAfterMs));
await sleep(backoffDelayMs(attempt, retryAfterMs), callerSignal);
attempt++;
continue;
}
Expand Down Expand Up @@ -174,15 +194,15 @@ export async function request<T>(
if (timedOut) {
const timeoutError = new RequestTimeoutError(timeoutMs);
if (attempt < retries) {
await sleep(backoffDelayMs(attempt));
await sleep(backoffDelayMs(attempt), callerSignal);
attempt++;
continue;
}
throw timeoutError;
}
if (isRetryableNetworkError(e) && attempt < retries) {
clearTimeout(timer);
await sleep(backoffDelayMs(attempt));
await sleep(backoffDelayMs(attempt), callerSignal);
attempt++;
continue;
}
Expand Down
51 changes: 35 additions & 16 deletions src/sandbox-client.ts
Original file line number Diff line number Diff line change
Expand Up @@ -263,24 +263,43 @@ export async function sandboxWait(
signal?: AbortSignal,
): Promise<SandboxDetail> {
const deadline = Date.now() + timeoutMs;
const deadlineController = new AbortController();
const abortFromCaller = () => deadlineController.abort();
const deadlineTimer = setTimeout(() => deadlineController.abort(), Math.max(0, timeoutMs));
if (signal?.aborted) deadlineController.abort();
else signal?.addEventListener('abort', abortFromCaller, { once: true });
let last: SandboxDetail | undefined;
while (Date.now() < deadline) {
if (signal?.aborted) throw new Error(`sandbox wait interrupted while waiting for ${wanted.join(' or ')}`);
last = await sandboxGet(opts, id, signal);
const state = String(last.observedState || '');
if (wanted.includes(state)) return last;
if (['FAILED', 'TERMINATED'].includes(state) && !wanted.includes(state)) {
throw new Error(`sandbox ${id} entered ${state} while waiting for ${wanted.join(' or ')}`);
try {
while (Date.now() < deadline) {
if (signal?.aborted) throw new Error(`sandbox wait interrupted while waiting for ${wanted.join(' or ')}`);
try {
last = await sandboxGet(opts, id, deadlineController.signal);
} catch (error) {
if (signal?.aborted) {
throw new Error(`sandbox wait interrupted while waiting for ${wanted.join(' or ')}`);
}
if (Date.now() >= deadline) break;
throw error;
}
if (Date.now() >= deadline) break;
const state = String(last.observedState || '');
if (wanted.includes(state)) return last;
if (['FAILED', 'TERMINATED'].includes(state) && !wanted.includes(state)) {
throw new Error(`sandbox ${id} entered ${state} while waiting for ${wanted.join(' or ')}`);
}
await new Promise<void>((resolve) => {
const done = () => {
clearTimeout(timer);
signal?.removeEventListener('abort', done);
resolve();
};
const timer = setTimeout(done, Math.min(intervalMs, Math.max(0, deadline - Date.now())));
signal?.addEventListener('abort', done, { once: true });
});
}
await new Promise<void>((resolve) => {
const done = () => {
clearTimeout(timer);
signal?.removeEventListener('abort', done);
resolve();
};
const timer = setTimeout(done, Math.min(intervalMs, Math.max(0, deadline - Date.now())));
signal?.addEventListener('abort', done, { once: true });
});
} finally {
clearTimeout(deadlineTimer);
signal?.removeEventListener('abort', abortFromCaller);
}
throw new Error(
`sandbox ${id} did not enter ${wanted.join(' or ')} within ${timeoutMs}ms` +
Expand Down
19 changes: 19 additions & 0 deletions src/tests/client.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -59,6 +59,25 @@ describe('client.request', () => {
expect(calls).toBe(1);
});

it('interrupts retry backoff when the caller aborts', async () => {
process.env.XAPI_RETRY_BASE_MS = '100';
let calls = 0;
fetchSpy = mockFetch(async () => {
calls++;
return new Response('busy', { status: 503 });
});
const controller = new AbortController();
const pending = request(
'https://action.xapi.to/x',
{ method: 'GET', signal: controller.signal },
5_000,
2,
);
controller.abort();
await expect(pending).rejects.toThrow(/aborted/i);
expect(calls).toBe(1);
});

it('does NOT retry by default (fail-safe for non-idempotent writes)', async () => {
let calls = 0;
fetchSpy = mockFetch(async () => {
Expand Down
39 changes: 39 additions & 0 deletions src/tests/sandbox-client.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -155,6 +155,45 @@ describe('sandbox client', () => {
expect(calls).toBe(3);
});

it('aborts an in-flight state read when the wait deadline expires', async () => {
let calls = 0;
fetchSpy = spyOn(globalThis, 'fetch').mockImplementation(((_url: any, init: any) => {
calls++;
return new Promise((_resolve, reject) => {
init?.signal?.addEventListener(
'abort',
() => reject(new DOMException('aborted', 'AbortError')),
{ once: true },
);
});
}) as any);
await expect(client.sandboxWait(
{ sandboxHost: 'sandbox.test.xapi.to', apiKey: 'sk-test' },
'box-1',
['RUNNING'],
20,
1,
)).rejects.toThrow(/within 20ms/);
expect(calls).toBe(1);
});

it('does not accept a desired state returned after the wait deadline', async () => {
fetchSpy = spyOn(globalThis, 'fetch').mockImplementation((async () => {
await new Promise((resolve) => setTimeout(resolve, 30));
return new Response(JSON.stringify({ id: 'box-1', observedState: 'RUNNING' }), {
status: 200,
headers: { 'content-type': 'application/json' },
});
}) as any);
await expect(client.sandboxWait(
{ sandboxHost: 'sandbox.test.xapi.to', apiKey: 'sk-test' },
'box-1',
['RUNNING'],
5,
1,
)).rejects.toThrow(/within 5ms/);
});

it('interrupts a state wait promptly so callers can clean up', async () => {
fetchSpy = spyOn(globalThis, 'fetch').mockResolvedValue(new Response(JSON.stringify({
id: 'box-1', observedState: 'PROVISIONING',
Expand Down