diff --git a/src/components/map-projects/MapProject.jsx b/src/components/map-projects/MapProject.jsx index 668a79b..74425f1 100644 --- a/src/components/map-projects/MapProject.jsx +++ b/src/components/map-projects/MapProject.jsx @@ -116,6 +116,7 @@ import { normalizeAlgorithmInvocation, hasSuccessfulAlgorithmResponse, getAlgori import { parseConceptKey } from './conceptKey' import { getDefaultTargetRepoVersion, getProjectTargetRepoVersion, getTargetRepoVersionFromUrl, getTargetRepoVersionId } from './projectTargetRepo' import { buildBridgeTargetDownloadEntries, buildQualityRowViews, conceptBelongsToTargetRepo, conceptForMapping, formatBridgeTargetDownloadEntry, resolveAICandidateID, getScoreDetails, getAIAnalysisCandidateIDs } from './viewBuilders.js' +import { getCapacityWaitLabel, mergeCapacityWaits } from './rowProgress.js' import './MapProject.scss' import '../common/ResizablePanel.scss' @@ -240,7 +241,7 @@ const MapProject = () => { // The requests waiting for capacity now, each with its rows, for the // "Waiting for capacity" notices. const capacityWaitsRef = React.useRef(new Map()) - const [capacityWaitRows, setCapacityWaitRows] = React.useState(null) + const [capacityWaits, setCapacityWaits] = React.useState(null) // A run's rows the server stayed too busy for, by algorithm id ('rerank' // included), for the end-of-run notice. const throttledRunRowsRef = React.useRef({}) @@ -1759,31 +1760,38 @@ const MapProject = () => { } const syncCapacityWaits = () => { + if(!capacityWaitsRef.current.size) + return setCapacityWaits(null) const rows = {} - capacityWaitsRef.current.forEach(rowIndexes => rowIndexes.forEach(index => { rows[index] = true })) - setCapacityWaitRows(capacityWaitsRef.current.size ? rows : null) + let all = null + capacityWaitsRef.current.forEach(({rowIndexes, wait}) => { + rowIndexes.forEach(index => { rows[index] = mergeCapacityWaits(rows[index], wait) }) + all = mergeCapacityWaits(all, wait) + }) + setCapacityWaits({rows, wait: all}) } // onWait/onWaitEnd for one request on these rows. While it waits out a busy // server (a 429, or a pause another request's 429 started; not an error - // backoff), its rows show "Waiting for capacity", and its first wait goes in + // backoff), its rows show why (ocl_issues#2865), and its first wait goes in // each row's log with the server's capacity headers. const trackCapacityWait = (rowIndexes, logExtras = {}) => { const waitId = {} let logged = false return { - onWait: ({reason, retryAfterMs, capacity}) => { + onWait: ({reason, delayMs, retryAfterMs, capacity, limit}) => { if(reason === 'error') return - capacityWaitsRef.current.set(waitId, rowIndexes) + const wait = {limit, retryAt: Date.now() + delayMs} + capacityWaitsRef.current.set(waitId, {rowIndexes, wait}) syncCapacityWaits() if(logged) return logged = true rowIndexes.forEach(index => log({ action: 'capacity_wait', - description: t('map_project.waiting_for_capacity'), - extras: {...logExtras, reason, retry_after_ms: retryAfterMs ?? null, ...(capacity ? {capacity} : {})} + description: getCapacityWaitLabel(wait, {t}), + extras: {...logExtras, reason, limit, retry_after_ms: retryAfterMs ?? null, ...(capacity ? {capacity} : {})} }, index)) }, onWaitEnd: () => { @@ -6049,14 +6057,14 @@ const MapProject = () => { // A request is waiting out a busy server: the run is // slower, not stuck (ocl_issues#2849). // Short in the split view, where the full notice would be cut off. - capacityWaitRows && - + capacityWaits && + } color='warning' variant='outlined' size='small' - label={isSplitView ? t('map_project.waiting_for_capacity_short') : t('map_project.waiting_for_capacity')} + label={getCapacityWaitLabel(capacityWaits.wait, {t, short: isSplitView})} sx={{margin: '5px'}} /> @@ -6466,7 +6474,7 @@ const MapProject = () => { candidatesScore={candidatesScore} rowIndex={rowIndex} rowStage={rowStageRef.current[rowIndex]} - capacityWait={Boolean(capacityWaitRows?.[rowIndex])} + capacityWait={capacityWaits?.rows?.[rowIndex] || false} rowState={rowMatchStateRef.current[rowIndex]} conceptCache={conceptCache} targetCanonical={buildProjectContext()?.target_repo?.canonical_url} diff --git a/src/components/map-projects/__tests__/rowProgress.test.js b/src/components/map-projects/__tests__/rowProgress.test.js index a94a37c..4775c97 100644 --- a/src/components/map-projects/__tests__/rowProgress.test.js +++ b/src/components/map-projects/__tests__/rowProgress.test.js @@ -10,7 +10,8 @@ import test from 'node:test' import assert from 'node:assert/strict' -import { getRowProgressLabel } from '../rowProgress.js' +import { getCapacityWaitLabel, getRowProgressLabel, mergeCapacityWaits } from '../rowProgress.js' +import { CAPACITY_LIMIT, RATE_LIMIT } from '../../../services/capacity.js' const ALGOS = [{id: 'ocl-semantic'}, {id: 'ocl-bridge'}] const t = key => key @@ -71,3 +72,42 @@ test('getRowProgressLabel: a done or failed rerank leaves the label as before', assert.equal(getRowProgressLabel({'ocl-semantic': 1, 'ocl-bridge': 1, rerank: 1}, ALGOS, {t}).label, undefined) assert.equal(getRowProgressLabel({'ocl-semantic': 1, 'ocl-bridge': 1, rerank: -2}, ALGOS, {t}).label, undefined) }) + +// ── ocl_issues#2865: capacity limit vs rate limit ─────────────────────────── + +const tWith = (key, values) => values ? `${key} ${JSON.stringify(values)}` : key +const formatTime = ms => `t+${ms}` +const capacityWait = {limit: CAPACITY_LIMIT, retryAt: 20000} +const rateLimitWait = {limit: RATE_LIMIT, retryAt: 20000} + +test('getCapacityWaitLabel: a capacity 429 keeps today\'s wording', () => { + assert.equal(getCapacityWaitLabel(capacityWait, {t: tWith, formatTime}), 'map_project.waiting_for_capacity') + assert.equal(getCapacityWaitLabel(capacityWait, {t: tWith, formatTime, short: true}), 'map_project.waiting_for_capacity_short') + assert.equal(getCapacityWaitLabel(true, {t: tWith, formatTime}), 'map_project.waiting_for_capacity') +}) + +test('getCapacityWaitLabel: a rate-limit 429 says the user is sending requests too quickly, and when it retries', () => { + assert.equal(getCapacityWaitLabel(rateLimitWait, {t: tWith, formatTime}), 'map_project.rate_limited {"time":"t+20000"}') + assert.equal(getCapacityWaitLabel(rateLimitWait, {t: tWith, formatTime, short: true}), 'map_project.rate_limited_short') +}) + +test('getRowProgressLabel: a row waiting out a rate-limit 429 says so, not "Waiting for capacity"', () => { + assert.deepEqual( + getRowProgressLabel({'ocl-semantic': 0, 'ocl-bridge': -1}, ALGOS, {t: tWith, capacityWait: rateLimitWait, formatTime}), + {label: 'map_project.rate_limited {"time":"t+20000"}', status: 'capacity_wait'}, + ) + assert.deepEqual( + getRowProgressLabel({'ocl-semantic': 0, 'ocl-bridge': -1}, ALGOS, {t: tWith, capacityWait, formatTime}), + {label: 'map_project.waiting_for_capacity', status: 'capacity_wait'}, + ) +}) + +test('mergeCapacityWaits: a capacity wait wins; between rate-limit waits, the later retry', () => { + const later = {limit: RATE_LIMIT, retryAt: 50000} + assert.equal(mergeCapacityWaits(null, rateLimitWait), rateLimitWait) + assert.equal(mergeCapacityWaits(rateLimitWait, undefined), rateLimitWait) + assert.equal(mergeCapacityWaits(rateLimitWait, capacityWait), capacityWait) + assert.equal(mergeCapacityWaits(capacityWait, later), capacityWait) + assert.equal(mergeCapacityWaits(rateLimitWait, later), later) + assert.equal(mergeCapacityWaits(later, rateLimitWait), later) +}) diff --git a/src/components/map-projects/rowProgress.js b/src/components/map-projects/rowProgress.js index 3f9d930..5022db7 100644 --- a/src/components/map-projects/rowProgress.js +++ b/src/components/map-projects/rowProgress.js @@ -1,3 +1,27 @@ +import { RATE_LIMIT } from '../../services/capacity.js' + +const formatClockTime = ms => new Date(ms).toLocaleTimeString([], {hour: 'numeric', minute: '2-digit', second: '2-digit'}) + +// wait is {limit, retryAt}; anything but a rate limit reads as a capacity wait (ocl_issues#2865). +export const getCapacityWaitLabel = (wait, { t, short = false, formatTime = formatClockTime } = {}) => { + if(wait?.limit === RATE_LIMIT) + return short ? + t('map_project.rate_limited_short') : + t('map_project.rate_limited', {time: formatTime(wait.retryAt)}) + return t(short ? 'map_project.waiting_for_capacity_short' : 'map_project.waiting_for_capacity') +} + +// A capacity wait wins; between two rate-limit waits, the later retry. +export const mergeCapacityWaits = (a, b) => { + if(!a || !b) + return a || b + if(a.limit !== RATE_LIMIT) + return a + if(b.limit !== RATE_LIMIT) + return b + return b.retryAt > a.retryAt ? b : a +} + /** * The row panel's progress chip: which of the row's algorithms is running or * still to run. {label: false} before the row has stages; no label once every @@ -8,11 +32,11 @@ * server stayed too busy for (-4) asks for a retry: it wasn't run, it didn't * fail. */ -export const getRowProgressLabel = (stageMap, algos, { t, capacityWait = false } = {}) => { +export const getRowProgressLabel = (stageMap, algos, { t, capacityWait = false, formatTime } = {}) => { if(stageMap === undefined) return {label: false} if(capacityWait) - return {label: t('map_project.waiting_for_capacity'), status: 'capacity_wait'} + return {label: getCapacityWaitLabel(capacityWait, {t, formatTime}), status: 'capacity_wait'} if(!stageMap) return {label: 'Preparing...', status: 'partial'} diff --git a/src/i18n/locales/en/translations.json b/src/i18n/locales/en/translations.json index c9595bc..86125ce 100644 --- a/src/i18n/locales/en/translations.json +++ b/src/i18n/locales/en/translations.json @@ -591,6 +591,8 @@ "match_request_failed": "Couldn't get candidates ({{error}}). Try again.", "waiting_for_capacity": "Waiting for capacity: results may take longer due to demand", "waiting_for_capacity_short": "Waiting for capacity", + "rate_limited": "You're sending requests too quickly: retrying at {{time}}", + "rate_limited_short": "Too many requests", "row_throttled": "Server busy: not matched yet. Retry later", "algorithm_throttled": "The server is busy, so this row wasn't matched this time. Nothing failed: run it again later.", "auto_match_rows_throttled": "The server was busy, so {{count}} input row(s) weren't matched: {{details}}. Nothing failed. Run Auto Match on those rows again later to finish them.", diff --git a/src/i18n/locales/es/translations.json b/src/i18n/locales/es/translations.json index f08cdad..669c8e3 100644 --- a/src/i18n/locales/es/translations.json +++ b/src/i18n/locales/es/translations.json @@ -555,6 +555,8 @@ "match_request_failed": "No se pudieron obtener candidatos ({{error}}). Inténtalo de nuevo.", "waiting_for_capacity": "Esperando capacidad: los resultados pueden tardar más debido a la demanda", "waiting_for_capacity_short": "Esperando capacidad", + "rate_limited": "Estás enviando solicitudes demasiado rápido: se reintentará a las {{time}}", + "rate_limited_short": "Demasiadas solicitudes", "row_throttled": "Servidor ocupado: aún sin procesar. Reintenta más tarde", "algorithm_throttled": "El servidor está ocupado, así que esta fila no se procesó esta vez. Nada falló: vuelve a ejecutarla más tarde.", "auto_match_rows_throttled": "El servidor estaba ocupado, así que {{count}} fila(s) de entrada no se procesaron: {{details}}. Nada falló. Vuelve a ejecutar Auto Match en esas filas más tarde para terminarlas.", diff --git a/src/i18n/locales/zh/translations.json b/src/i18n/locales/zh/translations.json index 0104747..3de810f 100644 --- a/src/i18n/locales/zh/translations.json +++ b/src/i18n/locales/zh/translations.json @@ -580,6 +580,8 @@ "match_request_failed": "无法获取候选项({{error}})。请重试。", "waiting_for_capacity": "正在等待容量:由于需求量大,结果可能需要更长时间", "waiting_for_capacity_short": "正在等待容量", + "rate_limited": "您发送请求的速度过快:将于 {{time}} 重试", + "rate_limited_short": "请求过多", "row_throttled": "服务器繁忙:尚未匹配。请稍后重试", "algorithm_throttled": "服务器繁忙,本次未能匹配此行。没有出错:请稍后重新运行。", "auto_match_rows_throttled": "服务器繁忙,{{count}} 个输入行未能匹配:{{details}}。没有出错。请稍后对这些行重新运行自动匹配以完成匹配。", diff --git a/src/services/__tests__/capacity.test.js b/src/services/__tests__/capacity.test.js index 062d8e5..99a4641 100644 --- a/src/services/__tests__/capacity.test.js +++ b/src/services/__tests__/capacity.test.js @@ -16,13 +16,16 @@ import test from 'node:test' import assert from 'node:assert/strict' import { + CAPACITY_LIMIT, CAPACITY_WAIT_CAP_MS, createCapacityGate, createLimiter, getCapacityHeaders, getRetryAfterMs, + getThrottleLimit, HEAVY_REQUEST_TIMEOUT_MS, LIGHT_REQUEST_TIMEOUT_MS, + RATE_LIMIT, requestWithCapacityRetry, sleepUnlessCancelled, untilCancelled, @@ -558,3 +561,92 @@ test('requestWithCapacityRetry: an attempt that timed out (axios ECONNABORTED) i assert.equal(result.ok, true) assert.equal(sent.length, 2) }) + +// ── ocl_issues#2865: capacity limit vs rate limit ─────────────────────────── + +const capacityExceeded = retryAfter => throttled(retryAfter, {data: {error_code: 'capacity_exceeded', detail: 'Server at capacity'}}) + +test('getThrottleLimit: capacity_exceeded is the capacity limit; any other 429 is the rate limit', () => { + assert.equal(getThrottleLimit(capacityExceeded(5).response), CAPACITY_LIMIT) + assert.equal(getThrottleLimit(throttled(5, {data: {detail: 'Request was throttled.'}}).response), RATE_LIMIT) + assert.equal(getThrottleLimit(throttled(5, {data: {error_code: 'rate_limited'}}).response), RATE_LIMIT) + assert.equal(getThrottleLimit(throttled(undefined, {data: 'Too Many Requests'}).response), RATE_LIMIT) + assert.equal(getThrottleLimit(undefined), RATE_LIMIT) +}) + +test('requestWithCapacityRetry: a capacity 429 waits its Retry-After, as the capacity limit', async () => { + const clock = virtualClock() + const waits = [] + const result = await requestWithCapacityRetry(scripted(capacityExceeded(20), ok()).send, { + now: clock.now, sleep: clock.sleep, random: () => 0, onWait: info => waits.push(info), + }) + + assert.equal(result.ok, true) + assert.equal(clock.t, 20000) + assert.deepEqual(waits.map(({reason, limit, delayMs, retryAfterMs}) => ({reason, limit, delayMs, retryAfterMs})), + [{reason: 'throttled', limit: CAPACITY_LIMIT, delayMs: 20000, retryAfterMs: 20000}]) +}) + +test('requestWithCapacityRetry: a rate-limit 429 with Retry-After waits it just the same, as the rate limit', async () => { + const clock = virtualClock() + const waits = [] + const result = await requestWithCapacityRetry(scripted(throttled(20, {data: {detail: 'Request was throttled.'}}), ok()).send, { + now: clock.now, sleep: clock.sleep, random: () => 0, onWait: info => waits.push(info), + }) + + assert.equal(result.ok, true) + assert.equal(clock.t, 20000) + assert.deepEqual(waits.map(({reason, limit, delayMs, retryAfterMs}) => ({reason, limit, delayMs, retryAfterMs})), + [{reason: 'throttled', limit: RATE_LIMIT, delayMs: 20000, retryAfterMs: 20000}]) +}) + +test('requestWithCapacityRetry: a rate-limit 429 without Retry-After backs off as before, as the rate limit, with its delay', async () => { + const clock = virtualClock() + const waits = [] + const result = await requestWithCapacityRetry(scripted(throttled(), throttled(), ok()).send, { + now: clock.now, sleep: clock.sleep, random: () => 0, onWait: info => waits.push(info), + }) + + assert.equal(result.ok, true) + assert.deepEqual(waits.map(({limit, delayMs, retryAfterMs}) => ({limit, delayMs, retryAfterMs})), [ + {limit: RATE_LIMIT, delayMs: 5000, retryAfterMs: null}, + {limit: RATE_LIMIT, delayMs: 10000, retryAfterMs: null}, + ]) +}) + +test('requestWithCapacityRetry: a request held by the shared gate gets the limit of the 429 that paused it', async () => { + const clock = virtualClock() + const heldBy = async limit => { + const gate = createCapacityGate({now: clock.now}) + gate.pause(10000, limit) + const waits = [] + await requestWithCapacityRetry(scripted(ok()).send, { + gate, now: clock.now, sleep: clock.sleep, random: () => 0, onWait: info => waits.push([info.reason, info.limit]), + }) + return waits + } + + assert.deepEqual(await heldBy(RATE_LIMIT), [['paused', RATE_LIMIT]]) + assert.deepEqual(await heldBy(CAPACITY_LIMIT), [['paused', CAPACITY_LIMIT]]) +}) + +test('createCapacityGate: the limit is that of the 429 holding the gate longest', () => { + const clock = virtualClock() + const gate = createCapacityGate({now: clock.now}) + gate.pause(30000, CAPACITY_LIMIT) + gate.pause(5000, RATE_LIMIT) + assert.equal(gate.limit(), CAPACITY_LIMIT) + gate.pause(60000, RATE_LIMIT) + assert.equal(gate.limit(), RATE_LIMIT) + gate.pause(120000) + assert.equal(gate.limit(), CAPACITY_LIMIT) +}) + +test('requestWithCapacityRetry: an error backoff carries no limit', async () => { + const clock = virtualClock() + const waits = [] + await requestWithCapacityRetry(scripted(httpError(503), ok()).send, { + now: clock.now, sleep: clock.sleep, random: () => 0, onWait: info => waits.push([info.reason, info.limit]), + }) + assert.deepEqual(waits, [['error', null]]) +}) diff --git a/src/services/capacity.js b/src/services/capacity.js index b4e0ab4..75184c0 100644 --- a/src/services/capacity.js +++ b/src/services/capacity.js @@ -77,6 +77,13 @@ export const getRetryAfterMs = (response, { now = Date.now } = {}) => { return typeof seconds === 'number' && Number.isFinite(seconds) && seconds >= 0 ? seconds * 1000 : null } +export const CAPACITY_LIMIT = 'capacity' +export const RATE_LIMIT = 'rate_limit' + +// Which limit refused a 429 (ocl_issues#2865). +export const getThrottleLimit = response => + response?.data?.error_code === 'capacity_exceeded' ? CAPACITY_LIMIT : RATE_LIMIT + const CAPACITY_HEADERS = { decision: ['x-ocl-capacity-decision', String], limit: ['x-ocl-capacity-limit', Number], @@ -130,10 +137,18 @@ export const sleepUnlessCancelled = async (ms, isCancelled = () => false, { slee */ export const createCapacityGate = ({ now = Date.now } = {}) => { let resumeAt = 0 + let pausedBy = CAPACITY_LIMIT const pausedForMs = () => Math.max(resumeAt - now(), 0) return { - pause: ms => { resumeAt = Math.max(resumeAt, now() + ms) }, + pause: (ms, limit = CAPACITY_LIMIT) => { + const at = now() + ms + if(at >= resumeAt) { + resumeAt = at + pausedBy = limit + } + }, pausedForMs, + limit: () => pausedBy, isPaused: () => pausedForMs() > 0, // Resolves true once the gate is open, false if cancelled first, or // 'timeout' once maxMs has passed with the gate still paused (other @@ -163,7 +178,7 @@ export const createCapacityGate = ({ now = Date.now } = {}) => { * @param {number} [opts.maxWaitMs] how long to wait in all before ending "throttled" * @param {number} [opts.maxRetries] retries for a network error or 502/503/504 (429s don't count) * @param {function} [opts.isRetryable] err => boolean; replaces the network-or-gateway check - * @param {function} [opts.onWait] ({reason: 'throttled'|'paused'|'error', delayMs, status, retryAfterMs, capacity}) => void + * @param {function} [opts.onWait] ({reason: 'throttled'|'paused'|'error', delayMs, status, retryAfterMs, capacity, limit}) => void * @param {function} [opts.onWaitEnd] () => void, after each wait, however it ended * @param {function} [opts.onCapacity] capacityHeaders => void, for each response that carries them * @returns {Promise<{ok: true, response, attempts, waitedMs}|{ok: false, reason: 'throttled'|'cancelled'|'error', error, attempts, waitedMs}>} @@ -215,7 +230,7 @@ export const requestWithCapacityRetry = async (send, { const budgetMs = maxWaitMs - waitedMs if(pauseMs > budgetMs) return end({ ok: false, reason: 'throttled', error: lastError }) - onWait?.({ reason: 'paused', delayMs: pauseMs, status: null, retryAfterMs: pauseMs }) + onWait?.({ reason: 'paused', delayMs: pauseMs, status: null, retryAfterMs: pauseMs, limit: gate.limit?.() ?? CAPACITY_LIMIT }) const startedAt = now() let outcome = await gate.wait(isCancelled, { sleep, pollMs, maxMs: budgetMs }) // Spread out the requests the pause held back, within the budget. @@ -241,6 +256,7 @@ export const requestWithCapacityRetry = async (send, { reportCapacity(err?.response) if(err?.response?.status === 429) { const retryAfterMs = getRetryAfterMs(err.response, { now }) + const limit = getThrottleLimit(err.response) const baseMs = retryAfterMs === null ? Math.min(DEFAULT_THROTTLE_WAIT_MS * 2 ** throttles, MAX_DEFAULT_THROTTLE_WAIT_MS) : Math.max(retryAfterMs, MIN_THROTTLE_WAIT_MS) @@ -252,8 +268,8 @@ export const requestWithCapacityRetry = async (send, { // asking for what was refused instead. if(waitedMs + delayMs > maxWaitMs) return end({ ok: false, reason: 'throttled', error: err }) - gate?.pause(baseMs) - if(!(await waitFor(delayMs, { reason: 'throttled', status: 429, retryAfterMs, capacity: getCapacityHeaders(err.response) }))) + gate?.pause(baseMs, limit) + if(!(await waitFor(delayMs, { reason: 'throttled', status: 429, retryAfterMs, capacity: getCapacityHeaders(err.response), limit }))) return cancelled() continue } @@ -263,7 +279,7 @@ export const requestWithCapacityRetry = async (send, { baseDelayMs * backoffFactor ** errorRetries * (1 - jitterFactor + random() * jitterFactor * 2) : Math.min(retryAfterMs, MAX_ERROR_RETRY_AFTER_MS) * (1 + random() * THROTTLE_JITTER) errorRetries += 1 - if(!(await waitFor(delayMs, { reason: 'error', status: err?.response?.status ?? null, retryAfterMs }))) + if(!(await waitFor(delayMs, { reason: 'error', status: err?.response?.status ?? null, retryAfterMs, limit: null }))) return cancelled() continue }