Skip to content

Commit 230de4d

Browse files
waleedlatif1claude
andauthored
fix(knowledge): wait for embedding admission instead of deferring documents for an hour (#7720)
* fix(knowledge): wait for embedding admission instead of deferring documents for an hour Every embedding batch on the indexing path waited at most five seconds for the deployment's shared admission bucket. Twenty concurrent documents fanning out eight batches each queue for minutes behind the configured per-minute budget, so under load most batches timed out, the document stopped, and it was re-dispatched with a delay that started at a minute and doubled to an hour. The bucket's own estimate of when capacity returns was only a floor under that ladder. During a bulk sync this produced thousands of hour-long deferrals for a limiter we run ourselves, while the provider was healthy. The knowledge path now waits up to a minute for admission, which is cheaper than the re-dispatch it replaces and still bounded by the per-request retry budget. The request bucket admits 64 concurrent starts instead of 8, so documents that begin together no longer lose a race for slots while the token budget sits unused. When an admission wait still expires, the document resumes after the bucket's stated wait, clamped to 10 to 60 seconds with jitter, and the yield counts against the processing-slice budget rather than the provider-failure attempts, since the provider did nothing wrong. The deadline path now carries the bucket's last stated wait so that estimate is available to the scheduler. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> * fix(knowledge): report only the admission wait still left at the deadline The bucket's stated wait is stored as an absolute instant so a deadline hit after a sleep carries the remainder, not the original duration. A test pins the knowledge admission wait below the retry budget the processing deadline reserves for each request. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> --------- Co-authored-by: Claude Fable 5.1 <noreply@anthropic.com>
1 parent 3878bd4 commit 230de4d

6 files changed

Lines changed: 154 additions & 11 deletions

File tree

‎apps/sim/lib/core/rate-limiter/provider-admission.test.ts‎

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -46,6 +46,8 @@ describe('provider admission', () => {
4646
'provider:embedding:openai:hashed-credential:requests',
4747
])
4848
expect(reservations[0].cost).toBe(50)
49+
/** Enough burst for every concurrent document to start a batch; the rate still governs throughput. */
50+
expect(reservations[1].config).toMatchObject({ maxTokens: 64, refillRate: 10 })
4951
expect(options.cooldownKeys).toHaveLength(2)
5052
}
5153
})
@@ -92,7 +94,7 @@ describe('provider admission', () => {
9294
})
9395
expect(consumeTokens).toHaveBeenCalledOnce()
9496
expect(consumeTokens.mock.calls[0][0]).toMatchObject([
95-
{ key: 'provider:ocr:openai:another-key:requests' },
97+
{ key: 'provider:ocr:openai:another-key:requests', config: { maxTokens: 2 } },
9698
])
9799
})
98100
it('retains the cooldown when an admission storage call consumes the remaining deadline', async () => {

‎apps/sim/lib/core/rate-limiter/provider-admission.ts‎

Lines changed: 23 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -18,6 +18,15 @@ interface ProviderAdmissionInput extends ProviderIdentity {
1818
maxWaitMs: number
1919
}
2020

21+
/**
22+
* Requests admitted in the same instant per embedding credential. The
23+
* per-minute rate still governs sustained throughput; the burst only decides
24+
* how many concurrent documents can start a batch together instead of losing a
25+
* race for a handful of slots while the token budget sits unused.
26+
*/
27+
const EMBEDDING_REQUEST_BURST = 64
28+
const DEFAULT_REQUEST_BURST = 2
29+
2130
/** A local admission wait expired; the document scheduler may retry the work later. */
2231
export class ProviderAdmissionTimeoutError extends Error {
2332
readonly retryable = false
@@ -68,15 +77,26 @@ export async function waitForProviderAdmission(input: ProviderAdmissionInput): P
6877
key: `${key}:requests`,
6978
cost: 1,
7079
config: {
71-
maxTokens: Math.min(input.operation === 'embedding' ? 8 : 2, requestsPerMinute),
80+
maxTokens: Math.min(
81+
input.operation === 'embedding' ? EMBEDDING_REQUEST_BURST : DEFAULT_REQUEST_BURST,
82+
requestsPerMinute
83+
),
7284
refillRate: requestsPerMinute / 60,
7385
refillIntervalMs: 1000,
7486
},
7587
})
7688

89+
/** When the bucket last said capacity returns, so a deadline hit after a sleep reports the wait still left. */
90+
let capacityAvailableAt: number | undefined
7791
for (;;) {
7892
input.signal?.throwIfAborted()
79-
if (Date.now() >= deadlineAt) throw new ProviderAdmissionTimeoutError()
93+
if (Date.now() >= deadlineAt) {
94+
const remainingMs =
95+
capacityAvailableAt === undefined ? undefined : capacityAvailableAt - Date.now()
96+
throw new ProviderAdmissionTimeoutError(
97+
remainingMs !== undefined && remainingMs > 0 ? remainingMs : undefined
98+
)
99+
}
80100
if (await isProviderQuotaExhausted(input))
81101
throw new ProviderQuotaExhaustedError(input.providerId)
82102
let result: AtomicAdmissionResult
@@ -98,6 +118,7 @@ export async function waitForProviderAdmission(input: ProviderAdmissionInput): P
98118
}
99119
if (result.allowed) return
100120
const waitMs = Math.max(1, result.retryAfterMs)
121+
if (Number.isFinite(waitMs)) capacityAvailableAt = Date.now() + waitMs
101122
if (!Number.isFinite(waitMs) || waitMs >= deadlineAt - Date.now()) {
102123
if (await isProviderQuotaExhausted(input))
103124
throw new ProviderQuotaExhaustedError(input.providerId)

‎apps/sim/lib/embeddings/client.test.ts‎

Lines changed: 12 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -10,6 +10,7 @@ import {
1010
assertKnowledgeEmbeddingCapacityForDeployment,
1111
clampEmbeddingConcurrency,
1212
EMBEDDING_MAX_RETRIES,
13+
EMBEDDING_RETRY_BUDGET_MS,
1314
EmbeddingAPIError,
1415
EmbeddingOutputLimitError,
1516
EmbeddingQuotaExhaustedError,
@@ -19,6 +20,7 @@ import {
1920
isBYOKEmbeddingCredentialRejection,
2021
isEmbeddingQuotaExhaustion,
2122
isTransientEmbeddingError,
23+
KNOWLEDGE_EMBEDDING_ADMISSION_WAIT_MS,
2224
MAX_EMBEDDING_SUCCESS_RESPONSE_BYTES,
2325
} from '@/lib/embeddings/client'
2426

@@ -1801,12 +1803,20 @@ describe('durable embedding batches', () => {
18011803
expect(fetchMock).not.toHaveBeenCalled()
18021804
})
18031805

1806+
it('keeps the checkpointed admission wait inside the retry budget the processing deadline reserves', () => {
1807+
expect(KNOWLEDGE_EMBEDDING_ADMISSION_WAIT_MS).toBeLessThan(EMBEDDING_RETRY_BUDGET_MS)
1808+
})
1809+
18041810
it('limits checkpointed admission waits while retaining the interactive request budget', async () => {
18051811
fetchMock.mockImplementation(() => Promise.resolve(jsonResponse(openAIBody([[1]], 7))))
18061812
await embed(['text'], { apiKey: 'fixture-key', checkpoints: memoryCheckpoints() })
1807-
expect(mockAdmit).toHaveBeenLastCalledWith(expect.objectContaining({ maxWaitMs: 5000 }))
1813+
expect(mockAdmit).toHaveBeenLastCalledWith(
1814+
expect.objectContaining({ maxWaitMs: KNOWLEDGE_EMBEDDING_ADMISSION_WAIT_MS })
1815+
)
18081816
await embed(['text'], { apiKey: 'fixture-key' })
1809-
expect(mockAdmit.mock.lastCall?.[0].maxWaitMs).toBeGreaterThan(5000)
1817+
expect(mockAdmit.mock.lastCall?.[0].maxWaitMs).toBeGreaterThan(
1818+
KNOWLEDGE_EMBEDDING_ADMISSION_WAIT_MS
1819+
)
18101820
})
18111821

18121822
it('drains admitted batches, resumes only missing requests and retains the complete token charge', async () => {

‎apps/sim/lib/embeddings/client.ts‎

Lines changed: 11 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -147,7 +147,17 @@ export const EMBEDDING_MAX_RETRY_DELAY_MS = 30_000
147147
* is honored in full when it fits inside this deadline.
148148
*/
149149
export const EMBEDDING_RETRY_BUDGET_MS = EMBEDDING_MAX_RETRIES * EMBEDDING_MAX_RETRY_DELAY_MS
150-
const KNOWLEDGE_EMBEDDING_ADMISSION_WAIT_MS = 5000
150+
151+
/**
152+
* How long a checkpointed indexing batch waits for the shared admission bucket
153+
* before the document yields its slot. Twenty concurrent documents fanning out
154+
* eight batches each can queue for a couple of minutes behind the configured
155+
* per-minute budget; yielding after a few seconds turned every such wait into a
156+
* full re-dispatch with a minute-or-more delay. A minute of idle waiting is far
157+
* cheaper than that round trip, and the per-request retry budget still bounds
158+
* the whole attempt. Interactive callers keep the full request budget.
159+
*/
160+
export const KNOWLEDGE_EMBEDDING_ADMISSION_WAIT_MS = 60_000
151161

152162
export class EmbeddingAPIError extends Error {
153163
public status: number

‎apps/sim/lib/knowledge/documents/processing-provider-continuation.test.ts‎

Lines changed: 67 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -15,6 +15,7 @@ import {
1515
MAX_PROCESSING_CONTINUATION_SLICES,
1616
MAX_PROVIDER_CONTINUATION_AGE_MS,
1717
MAX_PROVIDER_CONTINUATION_ATTEMPTS,
18+
resolveAdmissionContinuationDelayMs,
1819
resolveProviderContinuationDelayMs,
1920
scheduleDocumentProcessingProviderContinuation,
2021
} from '@/lib/knowledge/documents/processing-provider-continuation'
@@ -135,6 +136,55 @@ describe('durable provider continuations', () => {
135136
})
136137
})
137138

139+
it('resumes soon after the local admission bucket turned a batch away, without spending a provider attempt', async () => {
140+
const payload = {
141+
...PAYLOAD,
142+
providerRetryCount: 2,
143+
processingSliceCount: 7,
144+
providerRetryStartedAt: new Date(NOW.getTime() - 3_600_000).toISOString(),
145+
}
146+
const continuation = await scheduleDocumentProcessingProviderContinuation(
147+
payload,
148+
new ProviderCapacityDeferredError('admission_timeout', { retryAfterMs: 3_000 }),
149+
false
150+
)
151+
const delay = continuation.deferredUntil.getTime() - NOW.getTime()
152+
expect(delay).toBeGreaterThanOrEqual(8_000)
153+
expect(delay).toBeLessThanOrEqual(12_000)
154+
expect(continuation.processingQueueToken).toBe('knowledge-slice-doc-1-pass-1-8')
155+
expect(assertDocumentProcessingPayload(dispatch.mock.calls[0][0])).toMatchObject({
156+
providerRetryCount: 2,
157+
processingSliceCount: 8,
158+
providerRetryStartedAt: payload.providerRetryStartedAt,
159+
})
160+
})
161+
162+
it('keeps the exponential ladder for provider-side throttling', async () => {
163+
const continuation = await scheduleDocumentProcessingProviderContinuation(
164+
{ ...PAYLOAD, providerRetryCount: 3 },
165+
new ProviderCapacityDeferredError('rate_limit', { retryAfterMs: 3_000 }),
166+
false
167+
)
168+
const delay = continuation.deferredUntil.getTime() - NOW.getTime()
169+
expect(delay).toBeGreaterThanOrEqual(8 * 60_000 * 0.8)
170+
expect(continuation.processingQueueToken).toBe('knowledge-provider-doc-1-pass-1-4')
171+
})
172+
173+
it('bounds admission resumes independently of provider retries', async () => {
174+
await expect(
175+
scheduleDocumentProcessingProviderContinuation(
176+
{
177+
...PAYLOAD,
178+
processingSliceCount: MAX_PROCESSING_CONTINUATION_SLICES,
179+
providerRetryStartedAt: NOW.toISOString(),
180+
},
181+
new ProviderCapacityDeferredError('admission_timeout'),
182+
false
183+
)
184+
).rejects.toBeInstanceOf(ProviderCapacityContinuationExhaustedError)
185+
expect(dispatch).not.toHaveBeenCalled()
186+
})
187+
138188
it('starts the same bounded recovery horizon when the first continuation is a processing slice', async () => {
139189
await scheduleDocumentProcessingProviderContinuation(
140190
PAYLOAD,
@@ -187,6 +237,23 @@ describe('durable provider continuations', () => {
187237
).rejects.toBe(error)
188238
})
189239

240+
it('clamps and jitters the admission bucket wait', () => {
241+
for (let i = 0; i < 20; i++) {
242+
const stated = resolveAdmissionContinuationDelayMs(30_000)
243+
expect(stated).toBeGreaterThanOrEqual(24_000)
244+
expect(stated).toBeLessThanOrEqual(36_000)
245+
const floored = resolveAdmissionContinuationDelayMs(800)
246+
expect(floored).toBeGreaterThanOrEqual(8_000)
247+
expect(floored).toBeLessThanOrEqual(12_000)
248+
const capped = resolveAdmissionContinuationDelayMs(10 * 60_000)
249+
expect(capped).toBeGreaterThanOrEqual(48_000)
250+
expect(capped).toBeLessThanOrEqual(72_000)
251+
const missing = resolveAdmissionContinuationDelayMs(undefined)
252+
expect(missing).toBeGreaterThanOrEqual(12_000)
253+
expect(missing).toBeLessThanOrEqual(18_000)
254+
}
255+
})
256+
190257
it('bounds jittered polling without reducing provider minimums', () => {
191258
expect(resolveProviderContinuationDelayMs(1)).toBeGreaterThanOrEqual(48_000)
192259
expect(resolveProviderContinuationDelayMs(1)).toBeLessThanOrEqual(72_000)

‎apps/sim/lib/knowledge/documents/processing-provider-continuation.ts‎

Lines changed: 38 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -14,6 +14,16 @@ export const MAX_PROVIDER_CONTINUATION_ATTEMPTS = 48
1414
export const MAX_PROCESSING_CONTINUATION_SLICES = 512
1515
export const MAX_PROVIDER_CONTINUATION_AGE_MS = 24 * 60 * 60 * 1000
1616
const MAX_PROVIDER_CONTINUATION_DELAY_MS = 60 * 60 * 1000
17+
/**
18+
* Bounds for resuming after the deployment's own admission bucket ran out of
19+
* wait budget. The bucket states when capacity returns, but that estimate does
20+
* not know about the other documents waiting on it, so it is floored to keep
21+
* re-dispatches apart and capped so a document never sits idle for long while
22+
* the provider itself is healthy.
23+
*/
24+
const ADMISSION_RETRY_MIN_MS = 10_000
25+
const ADMISSION_RETRY_MAX_MS = 60_000
26+
const ADMISSION_RETRY_DEFAULT_MS = 15_000
1727

1828
/** Server-stated waits are a lower bound, including when they exceed the ordinary polling cap. */
1929
export function resolveProviderContinuationDelayMs(attempt: number, retryAfterMs?: number): number {
@@ -31,6 +41,20 @@ export function resolveProviderContinuationDelayMs(attempt: number, retryAfterMs
3141
)
3242
}
3343

44+
/**
45+
* Delay after the local admission bucket declined a batch: the bucket's stated
46+
* wait, clamped, with jitter so the documents it turned away do not return in
47+
* lockstep. Provider-side throttling keeps {@link resolveProviderContinuationDelayMs}.
48+
*/
49+
export function resolveAdmissionContinuationDelayMs(retryAfterMs?: number): number {
50+
const stated =
51+
retryAfterMs !== undefined && Number.isFinite(retryAfterMs) && retryAfterMs > 0
52+
? retryAfterMs
53+
: ADMISSION_RETRY_DEFAULT_MS
54+
const clamped = Math.min(Math.max(stated, ADMISSION_RETRY_MIN_MS), ADMISSION_RETRY_MAX_MS)
55+
return Math.round(backoffWithJitter(1, null, { baseMs: clamped, maxMs: clamped }))
56+
}
57+
3458
/** Defers capacity pressure without spending another document dispatch or changing billing identity. */
3559
export async function scheduleDocumentProcessingProviderContinuation(
3660
payload: DocumentProcessingPayload,
@@ -39,9 +63,16 @@ export async function scheduleDocumentProcessingProviderContinuation(
3963
predecessorAdmissionCharged = false
4064
): Promise<DocumentProcessingContinuation> {
4165
const now = Date.now()
66+
/**
67+
* A processing slice and a local admission timeout both mean the provider is
68+
* fine and the document simply needs another turn: neither spends one of the
69+
* bounded provider-failure attempts, and both count against the slice budget.
70+
*/
4271
const isProcessingSlice = error.reason === 'processing_budget'
43-
const providerRetryCount = (payload.providerRetryCount ?? 0) + (isProcessingSlice ? 0 : 1)
44-
const processingSliceCount = (payload.processingSliceCount ?? 0) + (isProcessingSlice ? 1 : 0)
72+
const isAdmissionTimeout = error.reason === 'admission_timeout'
73+
const isLocalYield = isProcessingSlice || isAdmissionTimeout
74+
const providerRetryCount = (payload.providerRetryCount ?? 0) + (isLocalYield ? 0 : 1)
75+
const processingSliceCount = (payload.processingSliceCount ?? 0) + (isLocalYield ? 1 : 0)
4576
const providerRetryStartedAt = payload.providerRetryStartedAt ?? new Date(now).toISOString()
4677
/** Tokenless legacy payloads retain a conservative handoff delay because their predecessor cannot be adopted safely. */
4778
const deferredUntil = new Date(
@@ -50,7 +81,9 @@ export async function scheduleDocumentProcessingProviderContinuation(
5081
? payload.processingQueueToken
5182
? 1000
5283
: 60_000
53-
: resolveProviderContinuationDelayMs(providerRetryCount, error.retryAfterMs))
84+
: isAdmissionTimeout
85+
? resolveAdmissionContinuationDelayMs(error.retryAfterMs)
86+
: resolveProviderContinuationDelayMs(providerRetryCount, error.retryAfterMs))
5487
)
5588
if (
5689
providerRetryCount > MAX_PROVIDER_CONTINUATION_ATTEMPTS ||
@@ -62,8 +95,8 @@ export async function scheduleDocumentProcessingProviderContinuation(
6295
}
6396
const processingQueueToken = createDocumentProcessingContinuationToken(
6497
payload,
65-
isProcessingSlice ? 'slice' : 'provider',
66-
isProcessingSlice ? processingSliceCount : providerRetryCount
98+
isLocalYield ? 'slice' : 'provider',
99+
isLocalYield ? processingSliceCount : providerRetryCount
67100
)
68101
await dispatchDocumentProcessingContinuation(
69102
{

0 commit comments

Comments
 (0)