Skip to content

Commit fe50edb

Browse files
waleedlatif1claude
andcommitted
fix(search): cap the bulk embedding lane below the credential budget
Every caller reserves from the credential's aggregate buckets, and the bulk lane additionally reserves from a bucket capped at 90% of that budget. The aggregate can no longer exceed the configured budget, and interactive callers always find headroom instead of a queue behind a crawl's batches. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01JBacX6HGVhPMUySuMfANwn
1 parent 82c9d39 commit fe50edb

4 files changed

Lines changed: 68 additions & 40 deletions

File tree

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

Lines changed: 25 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -86,17 +86,36 @@ describe('provider admission', () => {
8686
)
8787
})
8888

89-
it('reserves an interactive lane from its own buckets while sharing the provider gates', async () => {
89+
it('caps the bulk lane below the aggregate budget so interactive callers keep headroom', async () => {
90+
await waitForProviderAdmission({ ...INPUT, lane: 'bulk' })
9091
await waitForProviderAdmission({ ...INPUT, lane: 'interactive' })
91-
const [reservations, options] = consumeTokens.mock.calls[0]
92-
expect(reservations.map((item: { key: string }) => item.key)).toEqual([
93-
'provider:embedding:openai:hashed-credential:interactive:tokens',
94-
'provider:embedding:openai:hashed-credential:interactive:requests',
92+
const [bulkReservations, bulkOptions] = consumeTokens.mock.calls[0]
93+
expect(bulkReservations).toMatchObject([
94+
{
95+
key: 'provider:embedding:openai:hashed-credential:tokens',
96+
config: { maxTokens: 600_000, refillRate: 10_000 },
97+
},
98+
{ key: 'provider:embedding:openai:hashed-credential:requests', config: { maxTokens: 64 } },
99+
{
100+
key: 'provider:embedding:openai:hashed-credential:bulk:tokens',
101+
cost: 50,
102+
config: { maxTokens: 540_000, refillRate: 9_000 },
103+
},
104+
{
105+
key: 'provider:embedding:openai:hashed-credential:bulk:requests',
106+
config: { maxTokens: 57, refillRate: 9 },
107+
},
95108
])
96-
expect(options.cooldownKeys).toEqual([
109+
expect(bulkOptions.cooldownKeys).toEqual([
97110
'provider:embedding:openai:hashed-credential:cooldown',
98111
'provider:embedding:openai:hashed-credential:quota',
99112
])
113+
const [interactiveReservations, interactiveOptions] = consumeTokens.mock.calls[1]
114+
expect(interactiveReservations.map((item: { key: string }) => item.key)).toEqual([
115+
'provider:embedding:openai:hashed-credential:tokens',
116+
'provider:embedding:openai:hashed-credential:requests',
117+
])
118+
expect(interactiveOptions.cooldownKeys).toEqual(bulkOptions.cooldownKeys)
100119
})
101120

102121
it('isolates another credential and does not impose token costs on OCR', async () => {

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

Lines changed: 37 additions & 28 deletions
Original file line numberDiff line numberDiff line change
@@ -13,12 +13,16 @@ export interface ProviderIdentity {
1313
}
1414

1515
/**
16-
* A lane reserves from its own request and token buckets under the same
17-
* provider identity, so a person waiting on one embedding never queues behind
18-
* a bulk crawl's batches. Cooldown and quota gates stay per identity: a
19-
* provider pause or an exhausted balance still stops every lane.
16+
* Every caller reserves from the credential's aggregate buckets, so the
17+
* configured budget is never exceeded. The bulk lane also reserves from a
18+
* bucket capped at {@link BULK_LANE_SHARE} of that budget, which leaves an
19+
* interactive caller headroom instead of a queue behind a crawl's batches.
20+
* Cooldown and quota gates stay per identity: a provider pause or an exhausted
21+
* balance still stops every lane.
2022
*/
21-
export type ProviderAdmissionLane = 'interactive'
23+
export type ProviderAdmissionLane = 'bulk' | 'interactive'
24+
25+
const BULK_LANE_SHARE = 0.9
2226

2327
interface ProviderAdmissionInput extends ProviderIdentity {
2428
inputTokens?: number
@@ -58,43 +62,48 @@ export async function waitForProviderAdmission(input: ProviderAdmissionInput): P
5862
input.signal?.throwIfAborted()
5963
const deadlineAt = Date.now() + input.maxWaitMs
6064
const key = providerKey(input)
61-
const bucketKey = input.lane ? `${key}:${input.lane}` : key
6265
const requestsPerMinute =
6366
input.operation === 'embedding'
6467
? envNumber(env.KB_CONFIG_EMBEDDING_REQUESTS_PER_MINUTE, 600, { min: 1 })
6568
: input.operation === 'ocr'
6669
? envNumber(env.KB_CONFIG_OCR_REQUESTS_PER_MINUTE, 60, { min: 1 })
6770
: envNumber(env.KB_CONFIG_RERANK_REQUESTS_PER_MINUTE, 60, { min: 1 })
71+
const tokenBudget =
72+
input.operation === 'embedding' && input.inputTokens
73+
? {
74+
cost: input.inputTokens,
75+
perMinute: envNumber(env.KB_CONFIG_EMBEDDING_TOKENS_PER_MINUTE, 600_000, { min: 1 }),
76+
}
77+
: undefined
78+
if (tokenBudget && tokenBudget.cost > tokenBudget.perMinute) {
79+
throw new Error('Embedding request exceeds the configured per-credential token budget')
80+
}
81+
const requestBurst = Math.min(
82+
input.operation === 'embedding' ? EMBEDDING_REQUEST_BURST : DEFAULT_REQUEST_BURST,
83+
requestsPerMinute
84+
)
6885
const reservations: TokenBucketReservation[] = []
69-
if (input.operation === 'embedding' && input.inputTokens) {
70-
const tokensPerMinute = envNumber(env.KB_CONFIG_EMBEDDING_TOKENS_PER_MINUTE, 600_000, {
71-
min: 1,
72-
})
73-
if (input.inputTokens > tokensPerMinute) {
74-
throw new Error('Embedding request exceeds the configured per-credential token budget')
86+
const reserveBuckets = (bucketKey: string, share: number) => {
87+
if (tokenBudget) {
88+
const maxTokens = Math.floor(tokenBudget.perMinute * share)
89+
reservations.push({
90+
key: `${bucketKey}:tokens`,
91+
cost: tokenBudget.cost,
92+
config: { maxTokens, refillRate: maxTokens / 60, refillIntervalMs: 1000 },
93+
})
7594
}
7695
reservations.push({
77-
key: `${bucketKey}:tokens`,
78-
cost: input.inputTokens,
96+
key: `${bucketKey}:requests`,
97+
cost: 1,
7998
config: {
80-
maxTokens: tokensPerMinute,
81-
refillRate: tokensPerMinute / 60,
99+
maxTokens: Math.floor(requestBurst * share),
100+
refillRate: (requestsPerMinute * share) / 60,
82101
refillIntervalMs: 1000,
83102
},
84103
})
85104
}
86-
reservations.push({
87-
key: `${bucketKey}:requests`,
88-
cost: 1,
89-
config: {
90-
maxTokens: Math.min(
91-
input.operation === 'embedding' ? EMBEDDING_REQUEST_BURST : DEFAULT_REQUEST_BURST,
92-
requestsPerMinute
93-
),
94-
refillRate: requestsPerMinute / 60,
95-
refillIntervalMs: 1000,
96-
},
97-
})
105+
reserveBuckets(key, 1)
106+
if (input.lane === 'bulk') reserveBuckets(`${key}:bulk`, BULK_LANE_SHARE)
98107

99108
/** When the bucket last said capacity returns, so a deadline hit after a sleep reports the wait still left. */
100109
let capacityAvailableAt: number | undefined

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

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1807,11 +1807,11 @@ describe('durable embedding batches', () => {
18071807
expect(KNOWLEDGE_EMBEDDING_ADMISSION_WAIT_MS).toBeLessThan(EMBEDDING_RETRY_BUDGET_MS)
18081808
})
18091809

1810-
it('limits checkpointed admission waits and keeps interactive callers on their own lane', async () => {
1810+
it('limits checkpointed admission waits and keeps interactive callers off the bulk lane', async () => {
18111811
fetchMock.mockImplementation(() => Promise.resolve(jsonResponse(openAIBody([[1]], 7))))
18121812
await embed(['text'], { apiKey: 'fixture-key', checkpoints: memoryCheckpoints() })
18131813
expect(mockAdmit).toHaveBeenLastCalledWith(
1814-
expect.objectContaining({ maxWaitMs: KNOWLEDGE_EMBEDDING_ADMISSION_WAIT_MS, lane: undefined })
1814+
expect.objectContaining({ maxWaitMs: KNOWLEDGE_EMBEDDING_ADMISSION_WAIT_MS, lane: 'bulk' })
18151815
)
18161816
await embed(['text'], { apiKey: 'fixture-key' })
18171817
expect(mockAdmit).toHaveBeenLastCalledWith(expect.objectContaining({ lane: 'interactive' }))

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

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -853,11 +853,11 @@ async function callCheckpointedEmbeddingBatch(
853853
signal,
854854
checkpoints ? KNOWLEDGE_EMBEDDING_ADMISSION_WAIT_MS : undefined,
855855
/**
856-
* Checkpoints mark the bulk indexing path. Everything else has a person
857-
* waiting on it, so it reserves from the interactive lane and never
858-
* queues behind a crawl's batches.
856+
* Checkpoints mark the bulk indexing path, which may never take the whole
857+
* credential budget. Everything else has a person waiting on it and uses
858+
* the headroom the bulk lane leaves.
859859
*/
860-
checkpoints ? undefined : 'interactive'
860+
checkpoints ? 'bulk' : 'interactive'
861861
)
862862
if (identity) await checkpoints!.save(identity, result, signal)
863863
return result

0 commit comments

Comments
 (0)