mirror of
https://github.com/RooVetGit/Roo-Code.git
synced 2026-09-05 08:10:14 +00:00
feat: honor provider Retry-After headers on 429 responses in embedders
- Parse and honor Retry-After header from providers on rate limit errors - Support multiple header formats: Retry-After, X-RateLimit-Reset-After, X-RateLimit-Reset - Add support for Gemini structured retry info in error response body - Update global rate limit state to prefer provider-specified delays over exponential backoff - Add comprehensive tests for Retry-After header handling - Improve rate limit handling to reduce unnecessary delays and quota exhaustion Fixes #8101
This commit is contained in:
parent
87b45def18
commit
7b31ee7287
3 changed files with 598 additions and 37 deletions
|
|
@ -448,6 +448,265 @@ describe("OpenAICompatibleEmbedder", () => {
|
|||
})
|
||||
})
|
||||
|
||||
it("should honor Retry-After header when present", async () => {
|
||||
const testTexts = ["Hello world"]
|
||||
const testEmbedding = new Float32Array([0.25, 0.5, 0.75])
|
||||
const base64String = Buffer.from(testEmbedding.buffer).toString("base64")
|
||||
|
||||
// Create error with Retry-After header info
|
||||
const rateLimitError: any = {
|
||||
status: 429,
|
||||
message: "Rate limit exceeded",
|
||||
headers: {
|
||||
"retry-after": "3", // 3 seconds
|
||||
},
|
||||
rateLimitInfo: { retryAfterMs: 3000 },
|
||||
}
|
||||
|
||||
mockEmbeddingsCreate.mockRejectedValueOnce(rateLimitError).mockResolvedValueOnce({
|
||||
data: [{ embedding: base64String }],
|
||||
usage: { prompt_tokens: 10, total_tokens: 15 },
|
||||
})
|
||||
|
||||
const resultPromise = embedder.createEmbeddings(testTexts)
|
||||
|
||||
// First attempt fails immediately
|
||||
await vitest.advanceTimersByTimeAsync(100)
|
||||
|
||||
// Should wait for provider-specified 3 seconds (plus 1s buffer = 4s)
|
||||
await vitest.advanceTimersByTimeAsync(4000)
|
||||
|
||||
const result = await resultPromise
|
||||
|
||||
expect(mockEmbeddingsCreate).toHaveBeenCalledTimes(2)
|
||||
expect(console.warn).toHaveBeenCalledWith(expect.stringContaining("(using provider-specified delay)"))
|
||||
expect(result).toEqual({
|
||||
embeddings: [[0.25, 0.5, 0.75]],
|
||||
usage: { promptTokens: 10, totalTokens: 15 },
|
||||
})
|
||||
})
|
||||
|
||||
it("should parse Retry-After header as HTTP-date", async () => {
|
||||
const testTexts = ["Hello world"]
|
||||
const fullUrl = "https://api.example.com/v1/embeddings"
|
||||
const embedder = new OpenAICompatibleEmbedder(fullUrl, testApiKey, testModelId)
|
||||
|
||||
// Future date 5 seconds from now
|
||||
const futureDate = new Date(Date.now() + 5000)
|
||||
const httpDate = futureDate.toUTCString()
|
||||
|
||||
const mockFetch = global.fetch as MockedFunction<typeof fetch>
|
||||
mockFetch
|
||||
.mockResolvedValueOnce({
|
||||
ok: false,
|
||||
status: 429,
|
||||
headers: {
|
||||
get: (name: string) => (name === "retry-after" ? httpDate : null),
|
||||
},
|
||||
text: async () => "Rate limited",
|
||||
} as any)
|
||||
.mockResolvedValueOnce({
|
||||
ok: true,
|
||||
status: 200,
|
||||
json: async () => ({
|
||||
data: [{ embedding: [0.1, 0.2, 0.3] }],
|
||||
usage: { prompt_tokens: 10, total_tokens: 15 },
|
||||
}),
|
||||
} as any)
|
||||
|
||||
const resultPromise = embedder.createEmbeddings(testTexts)
|
||||
|
||||
// First attempt fails
|
||||
await vitest.advanceTimersByTimeAsync(100)
|
||||
|
||||
// Should wait approximately 5 seconds (plus buffer)
|
||||
await vitest.advanceTimersByTimeAsync(6000)
|
||||
|
||||
const result = await resultPromise
|
||||
|
||||
expect(mockFetch).toHaveBeenCalledTimes(2)
|
||||
expect(result.embeddings).toEqual([[0.1, 0.2, 0.3]])
|
||||
})
|
||||
|
||||
it("should handle X-RateLimit-Reset-After header", async () => {
|
||||
const testTexts = ["Hello world"]
|
||||
const fullUrl = "https://api.example.com/v1/embeddings"
|
||||
const embedder = new OpenAICompatibleEmbedder(fullUrl, testApiKey, testModelId)
|
||||
|
||||
const mockFetch = global.fetch as MockedFunction<typeof fetch>
|
||||
mockFetch
|
||||
.mockResolvedValueOnce({
|
||||
ok: false,
|
||||
status: 429,
|
||||
headers: {
|
||||
get: (name: string) => (name === "x-ratelimit-reset-after" ? "2" : null),
|
||||
},
|
||||
text: async () => "Rate limited",
|
||||
} as any)
|
||||
.mockResolvedValueOnce({
|
||||
ok: true,
|
||||
status: 200,
|
||||
json: async () => ({
|
||||
data: [{ embedding: [0.1, 0.2, 0.3] }],
|
||||
usage: { prompt_tokens: 10, total_tokens: 15 },
|
||||
}),
|
||||
} as any)
|
||||
|
||||
const resultPromise = embedder.createEmbeddings(testTexts)
|
||||
|
||||
// First attempt fails
|
||||
await vitest.advanceTimersByTimeAsync(100)
|
||||
|
||||
// Should wait 2 seconds (plus buffer)
|
||||
await vitest.advanceTimersByTimeAsync(3000)
|
||||
|
||||
const result = await resultPromise
|
||||
|
||||
expect(mockFetch).toHaveBeenCalledTimes(2)
|
||||
expect(result.embeddings).toEqual([[0.1, 0.2, 0.3]])
|
||||
})
|
||||
|
||||
it("should handle X-RateLimit-Reset header with Unix timestamp", async () => {
|
||||
const testTexts = ["Hello world"]
|
||||
const fullUrl = "https://api.example.com/v1/embeddings"
|
||||
const embedder = new OpenAICompatibleEmbedder(fullUrl, testApiKey, testModelId)
|
||||
|
||||
// Unix timestamp 4 seconds in the future
|
||||
const resetTimestamp = Math.floor((Date.now() + 4000) / 1000)
|
||||
|
||||
const mockFetch = global.fetch as MockedFunction<typeof fetch>
|
||||
mockFetch
|
||||
.mockResolvedValueOnce({
|
||||
ok: false,
|
||||
status: 429,
|
||||
headers: {
|
||||
get: (name: string) => (name === "x-ratelimit-reset" ? resetTimestamp.toString() : null),
|
||||
},
|
||||
text: async () => "Rate limited",
|
||||
} as any)
|
||||
.mockResolvedValueOnce({
|
||||
ok: true,
|
||||
status: 200,
|
||||
json: async () => ({
|
||||
data: [{ embedding: [0.1, 0.2, 0.3] }],
|
||||
usage: { prompt_tokens: 10, total_tokens: 15 },
|
||||
}),
|
||||
} as any)
|
||||
|
||||
const resultPromise = embedder.createEmbeddings(testTexts)
|
||||
|
||||
// First attempt fails
|
||||
await vitest.advanceTimersByTimeAsync(100)
|
||||
|
||||
// Should wait approximately 4 seconds (plus buffer)
|
||||
await vitest.advanceTimersByTimeAsync(5000)
|
||||
|
||||
const result = await resultPromise
|
||||
|
||||
expect(mockFetch).toHaveBeenCalledTimes(2)
|
||||
expect(result.embeddings).toEqual([[0.1, 0.2, 0.3]])
|
||||
})
|
||||
|
||||
it("should handle Gemini-style structured retry info in error body", async () => {
|
||||
const testTexts = ["Hello world"]
|
||||
const fullUrl = "https://generativelanguage.googleapis.com/v1beta/openai/embeddings"
|
||||
const embedder = new OpenAICompatibleEmbedder(fullUrl, testApiKey, testModelId)
|
||||
|
||||
const errorBody = {
|
||||
error: {
|
||||
code: 429,
|
||||
message: "Resource exhausted",
|
||||
details: [
|
||||
{
|
||||
metadata: {
|
||||
retry_delay: "10s",
|
||||
},
|
||||
},
|
||||
],
|
||||
},
|
||||
}
|
||||
|
||||
const mockFetch = global.fetch as MockedFunction<typeof fetch>
|
||||
mockFetch
|
||||
.mockResolvedValueOnce({
|
||||
ok: false,
|
||||
status: 429,
|
||||
headers: {
|
||||
get: () => null,
|
||||
},
|
||||
text: async () => JSON.stringify(errorBody),
|
||||
} as any)
|
||||
.mockResolvedValueOnce({
|
||||
ok: true,
|
||||
status: 200,
|
||||
json: async () => ({
|
||||
data: [{ embedding: [0.1, 0.2, 0.3] }],
|
||||
usage: { prompt_tokens: 10, total_tokens: 15 },
|
||||
}),
|
||||
} as any)
|
||||
|
||||
const resultPromise = embedder.createEmbeddings(testTexts)
|
||||
|
||||
// First attempt fails
|
||||
await vitest.advanceTimersByTimeAsync(100)
|
||||
|
||||
// Should wait 10 seconds (plus buffer)
|
||||
await vitest.advanceTimersByTimeAsync(11000)
|
||||
|
||||
const result = await resultPromise
|
||||
|
||||
expect(mockFetch).toHaveBeenCalledTimes(2)
|
||||
expect(result.embeddings).toEqual([[0.1, 0.2, 0.3]])
|
||||
})
|
||||
|
||||
it("should parse duration strings correctly", async () => {
|
||||
const embedder = new OpenAICompatibleEmbedder(testBaseUrl, testApiKey, testModelId)
|
||||
|
||||
// Access private method for testing
|
||||
const parseDurationString = (embedder as any).parseDurationString.bind(embedder)
|
||||
|
||||
expect(parseDurationString("10s")).toBe(10000)
|
||||
expect(parseDurationString("2m")).toBe(120000)
|
||||
expect(parseDurationString("1h")).toBe(3600000)
|
||||
expect(parseDurationString("invalid")).toBeUndefined()
|
||||
expect(parseDurationString(null)).toBeUndefined()
|
||||
expect(parseDurationString("")).toBeUndefined()
|
||||
})
|
||||
|
||||
it("should fall back to exponential backoff when no Retry-After is provided", async () => {
|
||||
const testTexts = ["Hello world"]
|
||||
const rateLimitError = {
|
||||
status: 429,
|
||||
message: "Rate limit exceeded",
|
||||
// No headers or rateLimitInfo
|
||||
}
|
||||
|
||||
const testEmbedding = new Float32Array([0.25, 0.5, 0.75])
|
||||
const base64String = Buffer.from(testEmbedding.buffer).toString("base64")
|
||||
|
||||
mockEmbeddingsCreate.mockRejectedValueOnce(rateLimitError).mockResolvedValueOnce({
|
||||
data: [{ embedding: base64String }],
|
||||
usage: { prompt_tokens: 10, total_tokens: 15 },
|
||||
})
|
||||
|
||||
const resultPromise = embedder.createEmbeddings(testTexts)
|
||||
|
||||
// First attempt fails
|
||||
await vitest.advanceTimersByTimeAsync(100)
|
||||
|
||||
// Should use exponential backoff (5s for first retry)
|
||||
await vitest.advanceTimersByTimeAsync(5000)
|
||||
|
||||
const result = await resultPromise
|
||||
|
||||
expect(mockEmbeddingsCreate).toHaveBeenCalledTimes(2)
|
||||
expect(console.warn).toHaveBeenCalledWith(expect.stringContaining("(using exponential backoff)"))
|
||||
expect(result).toEqual({
|
||||
embeddings: [[0.25, 0.5, 0.75]],
|
||||
usage: { promptTokens: 10, totalTokens: 15 },
|
||||
})
|
||||
})
|
||||
|
||||
it("should not retry on non-rate-limit errors", async () => {
|
||||
const testTexts = ["Hello world"]
|
||||
const authError = new Error("Unauthorized")
|
||||
|
|
|
|||
|
|
@ -27,6 +27,11 @@ interface OpenAIEmbeddingResponse {
|
|||
}
|
||||
}
|
||||
|
||||
interface RateLimitInfo {
|
||||
retryAfterMs?: number
|
||||
retryAfterDate?: Date
|
||||
}
|
||||
|
||||
/**
|
||||
* OpenAI Compatible implementation of the embedder interface with batching and rate limiting.
|
||||
* This embedder allows using any OpenAI-compatible API endpoint by specifying a custom baseURL.
|
||||
|
|
@ -191,6 +196,134 @@ export class OpenAICompatibleEmbedder implements IEmbedder {
|
|||
return patterns.some((pattern) => pattern.test(url))
|
||||
}
|
||||
|
||||
/**
|
||||
* Parses the Retry-After header to determine wait time
|
||||
* @param retryAfter The Retry-After header value
|
||||
* @returns The number of milliseconds to wait, or undefined if not parseable
|
||||
*/
|
||||
private parseRetryAfter(retryAfter: string | null): number | undefined {
|
||||
if (!retryAfter) return undefined
|
||||
|
||||
// Check if it's a delay in seconds (numeric value)
|
||||
const seconds = parseInt(retryAfter, 10)
|
||||
if (!isNaN(seconds)) {
|
||||
return seconds * 1000
|
||||
}
|
||||
|
||||
// Check if it's an HTTP-date
|
||||
const retryDate = new Date(retryAfter)
|
||||
if (!isNaN(retryDate.getTime())) {
|
||||
const now = Date.now()
|
||||
const delay = retryDate.getTime() - now
|
||||
return delay > 0 ? delay : 0
|
||||
}
|
||||
|
||||
return undefined
|
||||
}
|
||||
|
||||
/**
|
||||
* Extracts rate limit information from response headers and body
|
||||
* @param response The fetch Response object
|
||||
* @param errorBody Optional error response body for providers that include retry info there
|
||||
* @returns Rate limit information if available
|
||||
*/
|
||||
private extractRateLimitInfo(response: Response | null, errorBody?: any): RateLimitInfo {
|
||||
const info: RateLimitInfo = {}
|
||||
|
||||
if (!response) return info
|
||||
|
||||
// Helper function to safely get header value
|
||||
const getHeader = (name: string): string | null => {
|
||||
if (response.headers) {
|
||||
// Check if it's a proper Headers object with get method
|
||||
if (typeof response.headers.get === "function") {
|
||||
return response.headers.get(name)
|
||||
}
|
||||
// Check if it's a plain object
|
||||
if (typeof response.headers === "object") {
|
||||
return (response.headers as any)[name] || null
|
||||
}
|
||||
}
|
||||
return null
|
||||
}
|
||||
|
||||
// Standard Retry-After header (used by most providers)
|
||||
const retryAfter = getHeader("retry-after")
|
||||
if (retryAfter) {
|
||||
const delayMs = this.parseRetryAfter(retryAfter)
|
||||
if (delayMs !== undefined) {
|
||||
info.retryAfterMs = delayMs
|
||||
}
|
||||
}
|
||||
|
||||
// X-RateLimit-Reset-After header (used by some providers like Anthropic)
|
||||
const resetAfter = getHeader("x-ratelimit-reset-after")
|
||||
if (resetAfter) {
|
||||
const seconds = parseInt(resetAfter, 10)
|
||||
if (!isNaN(seconds)) {
|
||||
info.retryAfterMs = seconds * 1000
|
||||
}
|
||||
}
|
||||
|
||||
// X-RateLimit-Reset header (Unix timestamp)
|
||||
const resetTimestamp = getHeader("x-ratelimit-reset")
|
||||
if (resetTimestamp) {
|
||||
const timestamp = parseInt(resetTimestamp, 10)
|
||||
if (!isNaN(timestamp)) {
|
||||
const resetTime = timestamp * 1000 // Convert to milliseconds
|
||||
const now = Date.now()
|
||||
const delay = resetTime - now
|
||||
if (delay > 0) {
|
||||
info.retryAfterMs = delay
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Check for Gemini-specific retry information in error body
|
||||
if (errorBody && typeof errorBody === "object") {
|
||||
// Gemini may include retry information in the error response
|
||||
if (errorBody.error?.details) {
|
||||
for (const detail of errorBody.error.details) {
|
||||
if (detail.metadata?.retry_delay) {
|
||||
// Parse duration string like "10s" or "1m"
|
||||
const delay = this.parseDurationString(detail.metadata.retry_delay)
|
||||
if (delay) {
|
||||
info.retryAfterMs = delay
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
return info
|
||||
}
|
||||
|
||||
/**
|
||||
* Parses duration strings like "10s", "1m", "1h" to milliseconds
|
||||
* @param duration Duration string
|
||||
* @returns Duration in milliseconds or undefined
|
||||
*/
|
||||
private parseDurationString(duration: string): number | undefined {
|
||||
if (!duration || typeof duration !== "string") return undefined
|
||||
|
||||
const match = duration.match(/^(\d+)([smh])$/)
|
||||
if (!match) return undefined
|
||||
|
||||
const value = parseInt(match[1], 10)
|
||||
const unit = match[2]
|
||||
|
||||
switch (unit) {
|
||||
case "s":
|
||||
return value * 1000
|
||||
case "m":
|
||||
return value * 60 * 1000
|
||||
case "h":
|
||||
return value * 60 * 60 * 1000
|
||||
default:
|
||||
return undefined
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Makes a direct HTTP request to the embeddings endpoint
|
||||
* Used when the user provides a full endpoint URL (e.g., Azure OpenAI with query parameters)
|
||||
|
|
@ -223,9 +356,16 @@ export class OpenAICompatibleEmbedder implements IEmbedder {
|
|||
if (!response || !response.ok) {
|
||||
const status = response?.status || 0
|
||||
let errorText = "No response"
|
||||
let errorBody: any
|
||||
try {
|
||||
if (response && typeof response.text === "function") {
|
||||
errorText = await response.text()
|
||||
// Try to parse as JSON for structured error info
|
||||
try {
|
||||
errorBody = JSON.parse(errorText)
|
||||
} catch {
|
||||
// Not JSON, keep as text
|
||||
}
|
||||
} else if (response) {
|
||||
errorText = `Error ${status}`
|
||||
}
|
||||
|
|
@ -233,8 +373,14 @@ export class OpenAICompatibleEmbedder implements IEmbedder {
|
|||
// Ignore text parsing errors
|
||||
errorText = `Error ${status}`
|
||||
}
|
||||
const error = new Error(`HTTP ${status}: ${errorText}`) as HttpError
|
||||
const error = new Error(`HTTP ${status}: ${errorText}`) as HttpError & { rateLimitInfo?: RateLimitInfo }
|
||||
error.status = status || response?.status || 0
|
||||
|
||||
// Extract rate limit info if this is a 429 error
|
||||
if (status === 429) {
|
||||
error.rateLimitInfo = this.extractRateLimitInfo(response, errorBody)
|
||||
}
|
||||
|
||||
throw error
|
||||
}
|
||||
|
||||
|
|
@ -266,20 +412,38 @@ export class OpenAICompatibleEmbedder implements IEmbedder {
|
|||
|
||||
try {
|
||||
let response: OpenAIEmbeddingResponse
|
||||
let rateLimitInfo: RateLimitInfo | undefined
|
||||
|
||||
if (isFullUrl) {
|
||||
// Use direct HTTP request for full endpoint URLs
|
||||
response = await this.makeDirectEmbeddingRequest(this.baseUrl, batchTexts, model)
|
||||
} else {
|
||||
// Use OpenAI SDK for base URLs
|
||||
response = (await this.embeddingsClient.embeddings.create({
|
||||
input: batchTexts,
|
||||
model: model,
|
||||
// OpenAI package (as of v4.78.1) has a parsing issue that truncates embedding dimensions to 256
|
||||
// when processing numeric arrays, which breaks compatibility with models using larger dimensions.
|
||||
// By requesting base64 encoding, we bypass the package's parser and handle decoding ourselves.
|
||||
encoding_format: "base64",
|
||||
})) as OpenAIEmbeddingResponse
|
||||
try {
|
||||
response = (await this.embeddingsClient.embeddings.create({
|
||||
input: batchTexts,
|
||||
model: model,
|
||||
// OpenAI package (as of v4.78.1) has a parsing issue that truncates embedding dimensions to 256
|
||||
// when processing numeric arrays, which breaks compatibility with models using larger dimensions.
|
||||
// By requesting base64 encoding, we bypass the package's parser and handle decoding ourselves.
|
||||
encoding_format: "base64",
|
||||
})) as OpenAIEmbeddingResponse
|
||||
} catch (sdkError: any) {
|
||||
// Extract rate limit info from SDK errors
|
||||
if (sdkError?.status === 429) {
|
||||
// The OpenAI SDK may include headers in the error
|
||||
if (sdkError.headers) {
|
||||
const mockResponse = {
|
||||
headers: new Map(Object.entries(sdkError.headers)),
|
||||
} as any
|
||||
rateLimitInfo = this.extractRateLimitInfo(mockResponse, sdkError.error)
|
||||
}
|
||||
// Re-throw with rate limit info attached
|
||||
const error = sdkError as HttpError & { rateLimitInfo?: RateLimitInfo }
|
||||
error.rateLimitInfo = rateLimitInfo
|
||||
}
|
||||
throw sdkError
|
||||
}
|
||||
}
|
||||
|
||||
// Convert base64 embeddings to float32 arrays
|
||||
|
|
@ -310,7 +474,7 @@ export class OpenAICompatibleEmbedder implements IEmbedder {
|
|||
totalTokens: response.usage?.total_tokens || 0,
|
||||
},
|
||||
}
|
||||
} catch (error) {
|
||||
} catch (error: any) {
|
||||
// Capture telemetry before error is reformatted
|
||||
TelemetryService.instance.captureEvent(TelemetryEventName.CODE_INDEX_ERROR, {
|
||||
error: error instanceof Error ? error.message : String(error),
|
||||
|
|
@ -322,24 +486,39 @@ export class OpenAICompatibleEmbedder implements IEmbedder {
|
|||
const hasMoreAttempts = attempts < MAX_RETRIES - 1
|
||||
|
||||
// Check if it's a rate limit error
|
||||
const httpError = error as HttpError
|
||||
const httpError = error as HttpError & { rateLimitInfo?: RateLimitInfo }
|
||||
if (httpError?.status === 429) {
|
||||
// Update global rate limit state
|
||||
// Update global rate limit state with provider-specific retry info
|
||||
await this.updateGlobalRateLimitState(httpError)
|
||||
|
||||
if (hasMoreAttempts) {
|
||||
// Calculate delay based on global rate limit state
|
||||
const baseDelay = INITIAL_DELAY_MS * Math.pow(2, attempts)
|
||||
const globalDelay = await this.getGlobalRateLimitDelay()
|
||||
const delayMs = Math.max(baseDelay, globalDelay)
|
||||
// Calculate delay based on provider guidance or fallback
|
||||
let delayMs: number
|
||||
|
||||
if (httpError.rateLimitInfo?.retryAfterMs) {
|
||||
// Use provider-specified delay
|
||||
delayMs = httpError.rateLimitInfo.retryAfterMs
|
||||
console.warn(
|
||||
t("embeddings:rateLimitRetry", {
|
||||
delayMs,
|
||||
attempt: attempts + 1,
|
||||
maxRetries: MAX_RETRIES,
|
||||
}) + " (using provider-specified delay)",
|
||||
)
|
||||
} else {
|
||||
// Fallback to exponential backoff
|
||||
const baseDelay = INITIAL_DELAY_MS * Math.pow(2, attempts)
|
||||
const globalDelay = await this.getGlobalRateLimitDelay()
|
||||
delayMs = Math.max(baseDelay, globalDelay)
|
||||
console.warn(
|
||||
t("embeddings:rateLimitRetry", {
|
||||
delayMs,
|
||||
attempt: attempts + 1,
|
||||
maxRetries: MAX_RETRIES,
|
||||
}) + " (using exponential backoff)",
|
||||
)
|
||||
}
|
||||
|
||||
console.warn(
|
||||
t("embeddings:rateLimitRetry", {
|
||||
delayMs,
|
||||
attempt: attempts + 1,
|
||||
maxRetries: MAX_RETRIES,
|
||||
}),
|
||||
)
|
||||
await new Promise((resolve) => setTimeout(resolve, delayMs))
|
||||
continue
|
||||
}
|
||||
|
|
@ -445,7 +624,7 @@ export class OpenAICompatibleEmbedder implements IEmbedder {
|
|||
/**
|
||||
* Updates global rate limit state when a 429 error occurs
|
||||
*/
|
||||
private async updateGlobalRateLimitState(error: HttpError): Promise<void> {
|
||||
private async updateGlobalRateLimitState(error: HttpError & { rateLimitInfo?: RateLimitInfo }): Promise<void> {
|
||||
const release = await OpenAICompatibleEmbedder.globalRateLimitState.mutex.acquire()
|
||||
try {
|
||||
const state = OpenAICompatibleEmbedder.globalRateLimitState
|
||||
|
|
@ -461,14 +640,23 @@ export class OpenAICompatibleEmbedder implements IEmbedder {
|
|||
|
||||
state.lastRateLimitError = now
|
||||
|
||||
// Calculate exponential backoff based on consecutive errors
|
||||
const baseDelay = 5000 // 5 seconds base
|
||||
const maxDelay = 300000 // 5 minutes max
|
||||
const exponentialDelay = Math.min(baseDelay * Math.pow(2, state.consecutiveRateLimitErrors - 1), maxDelay)
|
||||
let delay: number
|
||||
|
||||
// Prefer provider-specified delay if available
|
||||
if (error.rateLimitInfo?.retryAfterMs) {
|
||||
delay = error.rateLimitInfo.retryAfterMs
|
||||
// Add a small buffer to avoid hitting the limit immediately
|
||||
delay = Math.min(delay + 1000, 300000) // Cap at 5 minutes
|
||||
} else {
|
||||
// Fallback to exponential backoff
|
||||
const baseDelay = 5000 // 5 seconds base
|
||||
const maxDelay = 300000 // 5 minutes max
|
||||
delay = Math.min(baseDelay * Math.pow(2, state.consecutiveRateLimitErrors - 1), maxDelay)
|
||||
}
|
||||
|
||||
// Set global rate limit
|
||||
state.isRateLimited = true
|
||||
state.rateLimitResetTime = now + exponentialDelay
|
||||
state.rateLimitResetTime = now + delay
|
||||
|
||||
// Silent rate limit activation - no logging to prevent flooding
|
||||
} finally {
|
||||
|
|
|
|||
|
|
@ -125,6 +125,102 @@ export class OpenAiEmbedder extends OpenAiNativeHandler implements IEmbedder {
|
|||
return { embeddings: allEmbeddings, usage }
|
||||
}
|
||||
|
||||
/**
|
||||
* Parses the Retry-After header to determine wait time
|
||||
* @param retryAfter The Retry-After header value
|
||||
* @returns The number of milliseconds to wait, or undefined if not parseable
|
||||
*/
|
||||
private parseRetryAfter(retryAfter: string | null): number | undefined {
|
||||
if (!retryAfter) return undefined
|
||||
|
||||
// Check if it's a delay in seconds (numeric value)
|
||||
const seconds = parseInt(retryAfter, 10)
|
||||
if (!isNaN(seconds)) {
|
||||
return seconds * 1000
|
||||
}
|
||||
|
||||
// Check if it's an HTTP-date
|
||||
const retryDate = new Date(retryAfter)
|
||||
if (!isNaN(retryDate.getTime())) {
|
||||
const now = Date.now()
|
||||
const delay = retryDate.getTime() - now
|
||||
return delay > 0 ? delay : 0
|
||||
}
|
||||
|
||||
return undefined
|
||||
}
|
||||
|
||||
/**
|
||||
* Extracts rate limit information from OpenAI SDK error
|
||||
* @param error The error from OpenAI SDK
|
||||
* @returns The number of milliseconds to wait, or undefined
|
||||
*/
|
||||
private extractRateLimitDelay(error: any): number | undefined {
|
||||
// OpenAI SDK may include headers in the error object
|
||||
if (error?.headers) {
|
||||
// Try Retry-After header first
|
||||
const retryAfter = error.headers["retry-after"] || error.headers["Retry-After"]
|
||||
if (retryAfter) {
|
||||
const delay = this.parseRetryAfter(retryAfter)
|
||||
if (delay !== undefined) return delay
|
||||
}
|
||||
|
||||
// Try X-RateLimit-Reset-After (seconds)
|
||||
const resetAfter = error.headers["x-ratelimit-reset-after"] || error.headers["X-RateLimit-Reset-After"]
|
||||
if (resetAfter) {
|
||||
const seconds = parseInt(resetAfter, 10)
|
||||
if (!isNaN(seconds)) {
|
||||
return seconds * 1000
|
||||
}
|
||||
}
|
||||
|
||||
// Try X-RateLimit-Reset (Unix timestamp)
|
||||
const resetTimestamp = error.headers["x-ratelimit-reset"] || error.headers["X-RateLimit-Reset"]
|
||||
if (resetTimestamp) {
|
||||
const timestamp = parseInt(resetTimestamp, 10)
|
||||
if (!isNaN(timestamp)) {
|
||||
const resetTime = timestamp * 1000 // Convert to milliseconds
|
||||
const now = Date.now()
|
||||
const delay = resetTime - now
|
||||
if (delay > 0) return delay
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Check if the error response includes retry information
|
||||
if (error?.response?.headers) {
|
||||
const headers = error.response.headers
|
||||
|
||||
// Try the same header checks on response.headers
|
||||
const retryAfter = headers.get?.("retry-after") || headers["retry-after"]
|
||||
if (retryAfter) {
|
||||
const delay = this.parseRetryAfter(retryAfter)
|
||||
if (delay !== undefined) return delay
|
||||
}
|
||||
|
||||
const resetAfter = headers.get?.("x-ratelimit-reset-after") || headers["x-ratelimit-reset-after"]
|
||||
if (resetAfter) {
|
||||
const seconds = parseInt(resetAfter, 10)
|
||||
if (!isNaN(seconds)) {
|
||||
return seconds * 1000
|
||||
}
|
||||
}
|
||||
|
||||
const resetTimestamp = headers.get?.("x-ratelimit-reset") || headers["x-ratelimit-reset"]
|
||||
if (resetTimestamp) {
|
||||
const timestamp = parseInt(resetTimestamp, 10)
|
||||
if (!isNaN(timestamp)) {
|
||||
const resetTime = timestamp * 1000
|
||||
const now = Date.now()
|
||||
const delay = resetTime - now
|
||||
if (delay > 0) return delay
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
return undefined
|
||||
}
|
||||
|
||||
/**
|
||||
* Helper method to handle batch embedding with retries and exponential backoff
|
||||
* @param batchTexts Array of texts to embed in this batch
|
||||
|
|
@ -155,14 +251,32 @@ export class OpenAiEmbedder extends OpenAiNativeHandler implements IEmbedder {
|
|||
// Check if it's a rate limit error
|
||||
const httpError = error as HttpError
|
||||
if (httpError?.status === 429 && hasMoreAttempts) {
|
||||
const delayMs = INITIAL_DELAY_MS * Math.pow(2, attempts)
|
||||
console.warn(
|
||||
t("embeddings:rateLimitRetry", {
|
||||
delayMs,
|
||||
attempt: attempts + 1,
|
||||
maxRetries: MAX_RETRIES,
|
||||
}),
|
||||
)
|
||||
// Try to extract provider-specified delay
|
||||
const providerDelay = this.extractRateLimitDelay(error)
|
||||
|
||||
let delayMs: number
|
||||
if (providerDelay !== undefined) {
|
||||
// Use provider-specified delay with a small buffer
|
||||
delayMs = Math.min(providerDelay + 1000, 300000) // Cap at 5 minutes
|
||||
console.warn(
|
||||
t("embeddings:rateLimitRetry", {
|
||||
delayMs,
|
||||
attempt: attempts + 1,
|
||||
maxRetries: MAX_RETRIES,
|
||||
}) + " (using provider-specified delay)",
|
||||
)
|
||||
} else {
|
||||
// Fallback to exponential backoff
|
||||
delayMs = INITIAL_DELAY_MS * Math.pow(2, attempts)
|
||||
console.warn(
|
||||
t("embeddings:rateLimitRetry", {
|
||||
delayMs,
|
||||
attempt: attempts + 1,
|
||||
maxRetries: MAX_RETRIES,
|
||||
}) + " (using exponential backoff)",
|
||||
)
|
||||
}
|
||||
|
||||
await new Promise((resolve) => setTimeout(resolve, delayMs))
|
||||
continue
|
||||
}
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue