Skip to content

Commit bba8f90

Browse files
waleedlatif1claude
andauthored
fix(search): reserve interactive embeddings from their own admission lane (#7726)
* fix(search): reserve interactive embeddings from their own admission lane A search embeds its query through the same per-credential admission bucket as bulk indexing, so during a crawl one short interactive request competed with hundreds of batches for the bucket's refill. Non-checkpointed callers now reserve from a separate lane under the same identity; the provider cooldown and quota gates stay shared, so a provider 429 or an exhausted balance still pauses every caller. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01JBacX6HGVhPMUySuMfANwn * 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 * fix(search): keep a bulk lane bucket large enough for one valid reservation At the smallest supported budgets the 90% share floored below a single request or token cost and locked the bulk lane. The bucket now holds at least the reservation it is asked for, while the refill rate keeps the share. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01JBacX6HGVhPMUySuMfANwn * refactor(search): derive the bulk admission lane from one flag A single bulk flag now selects the shorter admission wait and the capped lane, the lane cap is a pure function of configuration, and a bulk batch larger than that cap is rejected up front instead of resizing a shared bucket. Tests use the shared env mock. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> --------- Co-authored-by: Claude Fable 5.1 <noreply@anthropic.com>
1 parent 6d42e7d commit bba8f90

4 files changed

Lines changed: 98 additions & 27 deletions

File tree

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

Lines changed: 46 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,7 @@
11
/**
22
* @vitest-environment node
33
*/
4+
import { resetEnvMock, setEnv } from '@sim/testing/mocks/env.mock'
45
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
56

67
const { consumeTokens, getCooldownUntil, setCooldownUntil } = vi.hoisted(() => ({
@@ -35,7 +36,10 @@ describe('provider admission', () => {
3536
consumeTokens.mockResolvedValue({ allowed: true, tokensRemaining: 1, resetAt: new Date() })
3637
})
3738

38-
afterEach(() => vi.useRealTimers())
39+
afterEach(() => {
40+
vi.useRealTimers()
41+
resetEnvMock()
42+
})
3943

4044
it('shares both credential dimensions in one reservation across concurrent callers', async () => {
4145
await Promise.all([waitForProviderAdmission(INPUT), waitForProviderAdmission(INPUT)])
@@ -86,6 +90,47 @@ describe('provider admission', () => {
8690
)
8791
})
8892

93+
it('caps bulk work below the aggregate budget so interactive callers keep headroom', async () => {
94+
await waitForProviderAdmission({ ...INPUT, bulk: true })
95+
const [reservations, options] = consumeTokens.mock.calls[0]
96+
expect(reservations).toMatchObject([
97+
{ key: 'provider:embedding:openai:hashed-credential:tokens', config: { maxTokens: 600_000 } },
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+
},
108+
])
109+
expect(options.cooldownKeys).toEqual([
110+
'provider:embedding:openai:hashed-credential:cooldown',
111+
'provider:embedding:openai:hashed-credential:quota',
112+
])
113+
})
114+
115+
it('rejects a bulk batch the lane can never hold and keeps one request slot at a minimal burst', async () => {
116+
setEnv({
117+
KB_CONFIG_EMBEDDING_REQUESTS_PER_MINUTE: '1',
118+
KB_CONFIG_EMBEDDING_TOKENS_PER_MINUTE: '100',
119+
})
120+
await expect(
121+
waitForProviderAdmission({ ...INPUT, inputTokens: 95, bulk: true })
122+
).rejects.toThrow('exceeds the configured per-credential token budget')
123+
await waitForProviderAdmission({ ...INPUT, inputTokens: 95 })
124+
await waitForProviderAdmission({ ...INPUT, inputTokens: 90, bulk: true })
125+
expect(consumeTokens.mock.calls[1][0].slice(2)).toMatchObject([
126+
{ key: 'provider:embedding:openai:hashed-credential:bulk:tokens', config: { maxTokens: 90 } },
127+
{
128+
key: 'provider:embedding:openai:hashed-credential:bulk:requests',
129+
config: { maxTokens: 1 },
130+
},
131+
])
132+
})
133+
89134
it('isolates another credential and does not impose token costs on OCR', async () => {
90135
await waitForProviderAdmission({
91136
...INPUT,

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

Lines changed: 43 additions & 22 deletions
Original file line numberDiff line numberDiff line change
@@ -12,10 +12,20 @@ export interface ProviderIdentity {
1212
operation: 'embedding' | 'ocr' | 'rerank'
1313
}
1414

15+
/**
16+
* Share of a credential's budget the bulk lane may use. Every caller reserves
17+
* from the aggregate buckets, so the budget is never exceeded; bulk callers
18+
* also reserve from buckets capped at this share, which leaves an interactive
19+
* caller headroom instead of a queue behind a crawl's batches.
20+
*/
21+
const BULK_LANE_SHARE = 0.9
22+
1523
interface ProviderAdmissionInput extends ProviderIdentity {
1624
inputTokens?: number
1725
signal?: AbortSignal
1826
maxWaitMs: number
27+
/** Bulk work is capped at {@link BULK_LANE_SHARE}; cooldown and quota gates still stop every caller. */
28+
bulk?: boolean
1929
}
2030

2131
/**
@@ -55,36 +65,47 @@ export async function waitForProviderAdmission(input: ProviderAdmissionInput): P
5565
: input.operation === 'ocr'
5666
? envNumber(env.KB_CONFIG_OCR_REQUESTS_PER_MINUTE, 60, { min: 1 })
5767
: envNumber(env.KB_CONFIG_RERANK_REQUESTS_PER_MINUTE, 60, { min: 1 })
68+
const tokenBudget =
69+
input.operation === 'embedding' && input.inputTokens
70+
? {
71+
cost: input.inputTokens,
72+
perMinute: envNumber(env.KB_CONFIG_EMBEDDING_TOKENS_PER_MINUTE, 600_000, { min: 1 }),
73+
}
74+
: undefined
75+
const laneShare = input.bulk ? BULK_LANE_SHARE : 1
76+
if (tokenBudget && tokenBudget.cost > Math.floor(tokenBudget.perMinute * laneShare)) {
77+
throw new Error('Embedding request exceeds the configured per-credential token budget')
78+
}
79+
const requestBurst = Math.min(
80+
input.operation === 'embedding' ? EMBEDDING_REQUEST_BURST : DEFAULT_REQUEST_BURST,
81+
requestsPerMinute
82+
)
5883
const reservations: TokenBucketReservation[] = []
59-
if (input.operation === 'embedding' && input.inputTokens) {
60-
const tokensPerMinute = envNumber(env.KB_CONFIG_EMBEDDING_TOKENS_PER_MINUTE, 600_000, {
61-
min: 1,
62-
})
63-
if (input.inputTokens > tokensPerMinute) {
64-
throw new Error('Embedding request exceeds the configured per-credential token budget')
84+
const reserveBuckets = (bucketKey: string, share: number) => {
85+
if (tokenBudget) {
86+
reservations.push({
87+
key: `${bucketKey}:tokens`,
88+
cost: tokenBudget.cost,
89+
config: {
90+
maxTokens: Math.floor(tokenBudget.perMinute * share),
91+
refillRate: (tokenBudget.perMinute * share) / 60,
92+
refillIntervalMs: 1000,
93+
},
94+
})
6595
}
6696
reservations.push({
67-
key: `${key}:tokens`,
68-
cost: input.inputTokens,
97+
key: `${bucketKey}:requests`,
98+
cost: 1,
6999
config: {
70-
maxTokens: tokensPerMinute,
71-
refillRate: tokensPerMinute / 60,
100+
/** A burst of one leaves no share to carve out, so the lane then matches the aggregate. */
101+
maxTokens: Math.max(1, Math.floor(requestBurst * share)),
102+
refillRate: (requestsPerMinute * share) / 60,
72103
refillIntervalMs: 1000,
73104
},
74105
})
75106
}
76-
reservations.push({
77-
key: `${key}:requests`,
78-
cost: 1,
79-
config: {
80-
maxTokens: Math.min(
81-
input.operation === 'embedding' ? EMBEDDING_REQUEST_BURST : DEFAULT_REQUEST_BURST,
82-
requestsPerMinute
83-
),
84-
refillRate: requestsPerMinute / 60,
85-
refillIntervalMs: 1000,
86-
},
87-
})
107+
reserveBuckets(key, 1)
108+
if (input.bulk) reserveBuckets(`${key}:bulk`, BULK_LANE_SHARE)
88109

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

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

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

1810-
it('limits checkpointed admission waits while retaining the interactive request budget', 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 })
1814+
expect.objectContaining({ maxWaitMs: KNOWLEDGE_EMBEDDING_ADMISSION_WAIT_MS, bulk: true })
18151815
)
18161816
await embed(['text'], { apiKey: 'fixture-key' })
1817+
expect(mockAdmit).toHaveBeenLastCalledWith(expect.objectContaining({ bulk: false }))
18171818
expect(mockAdmit.mock.lastCall?.[0].maxWaitMs).toBeGreaterThan(
18181819
KNOWLEDGE_EMBEDDING_ADMISSION_WAIT_MS
18191820
)

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

Lines changed: 6 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -544,8 +544,10 @@ async function callEmbeddingAPI(
544544
expectedDimensions: number | undefined,
545545
isBYOK: boolean,
546546
signal?: AbortSignal,
547-
admissionWaitMs = EMBEDDING_RETRY_BUDGET_MS
547+
/** Bulk indexing waits briefly and is capped below the credential budget; everything else has a person waiting on it. */
548+
bulk = false
548549
): Promise<{ embeddings: number[][]; totalTokens: number; dimensions: number }> {
550+
const admissionWaitMs = bulk ? KNOWLEDGE_EMBEDDING_ADMISSION_WAIT_MS : EMBEDDING_RETRY_BUDGET_MS
549551
const admissionIdentity = embeddingAdmissionIdentity({ providerId, quotaCircuitIdentity, isBYOK })
550552
return retryWithExponentialBackoff(
551553
async (operationSignal, deadlineAt) => {
@@ -563,6 +565,7 @@ async function callEmbeddingAPI(
563565
),
564566
signal: operationSignal,
565567
maxWaitMs: Math.min(admissionWaitMs, Math.max(0, deadlineAt - Date.now())),
568+
bulk,
566569
})
567570
} catch (error) {
568571
if (error instanceof ProviderQuotaExhaustedError)
@@ -795,6 +798,7 @@ async function mapEmbeddingBatches<T, R>(
795798
return results.map((result) => result!.value)
796799
}
797800

801+
/** Checkpoints mark the bulk indexing path; every other caller is interactive. */
798802
async function callCheckpointedEmbeddingBatch(
799803
batch: string[],
800804
batchIndex: number,
@@ -848,7 +852,7 @@ async function callCheckpointedEmbeddingBatch(
848852
provider.dimensions,
849853
provider.isBYOK,
850854
signal,
851-
checkpoints ? KNOWLEDGE_EMBEDDING_ADMISSION_WAIT_MS : undefined
855+
checkpoints !== undefined
852856
)
853857
if (identity) await checkpoints!.save(identity, result, signal)
854858
return result

0 commit comments

Comments
 (0)