Skip to content
Open
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
4 changes: 2 additions & 2 deletions src/api/api.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1404,7 +1404,7 @@ export class OpenSeaAPI {
const raw = await this._fetch(url, "GET", undefined, undefined, options)
return this.camelizeResponseBody<T>(raw, options)
},
{ logger: this.logger },
{ logger: this.logger, signal: options?.signal },
)
}

Expand Down Expand Up @@ -1449,7 +1449,7 @@ export class OpenSeaAPI {
const raw = await this._fetch(url, method, headers, wireBody, options)
return this.camelizeResponseBody<T>(raw, options)
},
{ logger: this.logger },
{ logger: this.logger, signal: options?.signal },
)
}

Expand Down
32 changes: 29 additions & 3 deletions src/utils/rateLimit.ts
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,11 @@ const CUSTOM_RATE_LIMIT_STATUS_CODE = 599
export interface RateLimitOptions {
/** Logger function for logging progress */
logger?: (message: string) => void
/**
* Abort signal for cancelling a pending retry delay. The operation itself
* must handle cancellation separately.
*/
signal?: AbortSignal
/** Maximum number of retry attempts for rate limit errors */
maxRetries?: number
/** Base delay in ms to wait after a rate limit error if retry-after header is not present */
Expand All @@ -32,9 +37,29 @@ export interface RateLimitOptions {
/**
* Sleep for a specified duration
* @param ms Duration in milliseconds
* @param signal Optional signal for cancelling the delay
*/
function sleep(ms: number): Promise<void> {
return new Promise(resolve => setTimeout(resolve, ms))
function sleep(ms: number, signal?: AbortSignal): Promise<void> {
if (!signal) {
return new Promise(resolve => setTimeout(resolve, ms))
}
if (signal.aborted) {
return Promise.reject(new Error("Request aborted"))
}

return new Promise((resolve, reject) => {
let timeoutId: ReturnType<typeof setTimeout>
const onAbort = () => {
clearTimeout(timeoutId)
signal.removeEventListener("abort", onAbort)
reject(new Error("Request aborted"))
}
timeoutId = setTimeout(() => {
signal.removeEventListener("abort", onAbort)
resolve()
}, ms)
signal.addEventListener("abort", onAbort, { once: true })
})
}

/**
Expand All @@ -52,6 +77,7 @@ export async function executeWithRateLimit<T>(
): Promise<T> {
const {
logger = () => {},
signal,
maxRetries = DEFAULT_MAX_RETRIES,
baseRetryDelay = DEFAULT_BASE_RETRY_DELAY_MS,
} = options
Expand Down Expand Up @@ -92,7 +118,7 @@ export async function executeWithRateLimit<T>(
)
}

await sleep(delayMs)
await sleep(delayMs, signal)
logger(`Retrying operation...`)
}
}
Expand Down
48 changes: 48 additions & 0 deletions test/api/api.spec.ts
Original file line number Diff line number Diff line change
Expand Up @@ -204,6 +204,54 @@ describe("API", () => {
expect(result.address).toBe("0x0000000000000000000000000000000000000000")
})

test.each([
"get",
"request",
] as const)("API %s aborts while waiting to retry a rate-limited request", async method => {
const rateLimitError = Object.assign(new Error("429 Too Many Requests"), {
statusCode: 429,
retryAfter: 30,
}) as OpenSeaRateLimitError
const localApi = new OpenSeaAPI({ apiKey: "key" })
fetchStub = vi
.spyOn(
localApi as unknown as { _fetch: () => Promise<unknown> },
"_fetch",
)
.mockRejectedValueOnce(rateLimitError)
.mockResolvedValueOnce({ ok: true })
const controller = new AbortController()
let outcome: "pending" | "resolved" | "rejected" = "pending"
let rejection: unknown

const requestPromise =
method === "get"
? localApi.get("/api/v2/test", {}, { signal: controller.signal })
: localApi.request("POST", "/api/v2/test", undefined, undefined, {
signal: controller.signal,
})
const observed = requestPromise.then(
() => {
outcome = "resolved"
},
error => {
outcome = "rejected"
rejection = error
},
)

await vi.advanceTimersByTimeAsync(0)
expect(fetchStub).toHaveBeenCalledTimes(1)

controller.abort()
await vi.advanceTimersByTimeAsync(0)

expect(outcome).toBe("rejected")
expect((rejection as Error).message).toBe("Request aborted")
expect(fetchStub).toHaveBeenCalledTimes(1)
await observed
})

test("API parses Retry-After HTTP-date header", () => {
// Pin system time for deterministic date math
vi.setSystemTime(new Date("2020-01-01T00:00:00.000Z"))
Expand Down
39 changes: 39 additions & 0 deletions test/utils/rateLimit.spec.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@ import type { OpenSeaRateLimitError } from "../../src/types"
import {
executeSequentialWithRateLimit,
executeWithRateLimit,
type RateLimitOptions,
} from "../../src/utils/rateLimit"

describe("Utils: rateLimit", () => {
Expand Down Expand Up @@ -147,6 +148,44 @@ describe("Utils: rateLimit", () => {
expect(result).toBe("success")
expect(operation).toHaveBeenCalledTimes(2)
})

test("aborts while waiting to retry a rate-limited operation", async () => {
const rateLimitError = Object.assign(new Error("429 Too Many Requests"), {
statusCode: 429,
retryAfter: 30,
}) as OpenSeaRateLimitError
const operation = vi
.fn()
.mockRejectedValueOnce(rateLimitError)
.mockResolvedValueOnce("success")
const controller = new AbortController()
const options: RateLimitOptions = {
signal: controller.signal,
}
let outcome: "pending" | "resolved" | "rejected" = "pending"
let rejection: unknown

const observed = executeWithRateLimit(operation, options).then(
() => {
outcome = "resolved"
},
error => {
outcome = "rejected"
rejection = error
},
)

await vi.advanceTimersByTimeAsync(0)
expect(operation).toHaveBeenCalledTimes(1)

controller.abort()
await vi.advanceTimersByTimeAsync(0)

expect(outcome).toBe("rejected")
expect((rejection as Error).message).toBe("Request aborted")
expect(operation).toHaveBeenCalledTimes(1)
await observed
})
})

describe("executeSequentialWithRateLimit", () => {
Expand Down