diff --git a/src/components/map-projects/Candidates.jsx b/src/components/map-projects/Candidates.jsx
index 696ce03..cb2320e 100644
--- a/src/components/map-projects/Candidates.jsx
+++ b/src/components/map-projects/Candidates.jsx
@@ -23,6 +23,7 @@ import RefreshIcon from '@mui/icons-material/Refresh';
import SortIcon from '@mui/icons-material/Sort';
import GroupIcon from '@mui/icons-material/Layers';
import AssistantIcon from '@mui/icons-material/Assistant';
+import HourglassIcon from '@mui/icons-material/HourglassEmpty';
import isEmpty from 'lodash/isEmpty'
import flatten from 'lodash/flatten'
@@ -53,38 +54,7 @@ import {
conceptForMapping,
getAIAnalysisCandidateIDs
} from './viewBuilders.js'
-
-const getRowProgressLabel = (stageMap, algos) => {
- if(stageMap === undefined)
- return {label: false}
- if(!stageMap)
- return {label: 'Preparing...', status: 'partial'}
-
- const stages = algos.map(k => stageMap[k.id]);
-
- if (!stages.length || stages.every(v => v === -1)) {
- return { label: 'Not started', status: 'idle' };
- }
-
- const runningIndex = stages.findIndex(v => v === 0);
- if (runningIndex !== -1) {
- return {
- label: `Running: ${algos[runningIndex].id}...`,
- status: 'running',
- };
- }
-
- if (stages.every(v => v === 1)) {
- return true
- }
-
- const waitingIndex = stages.findIndex(v => v === -1)
- if(waitingIndex !== -1) {
- return {label: `Waiting: ${algos[waitingIndex].id}`, status: 'waiting'} // AutoMatch Bulk
- }
-
- return { label: 'Partially completed', status: 'partial' };
-}
+import { getRowProgressLabel } from './rowProgress.js'
const Sort = ({ selected, onSort }) => {
const { t } = useTranslation();
@@ -407,7 +377,7 @@ const CandidateList = ({rowViews, header, rowIndex, sortBy, order, openConceptPa
// conceptCache — project-wide ConceptDefinition store, keyed by concept_key.
// algosSelected — algorithm definitions (for headers/grouping).
// (plans/unified-mapper-model.md "How the views map onto this model".)
-const Candidates = ({rowIndex, rowState, conceptCache, targetCanonical, targetRelativeUrl, openConceptPanel, showItem, isSelectedForMap, onMap, onFetchMore, isLoading, candidatesScore, repoVersion, analysis, onFetchRecommendation, appliedFacets, setAppliedFacets, filters, facets, columns, defaultFilters, locales, models, selectedModel, onModelChange, promptTemplates, promptTemplate, onPromptTemplateChange, onRefreshClick, rowStage, inAIAssistantGroup, algosSelected, canSelectAIModel}) => {
+const Candidates = ({rowIndex, rowState, conceptCache, targetCanonical, targetRelativeUrl, openConceptPanel, showItem, isSelectedForMap, onMap, onFetchMore, isLoading, candidatesScore, repoVersion, analysis, onFetchRecommendation, appliedFacets, setAppliedFacets, filters, facets, columns, defaultFilters, locales, models, selectedModel, onModelChange, promptTemplates, promptTemplate, onPromptTemplateChange, onRefreshClick, rowStage, inAIAssistantGroup, algosSelected, canSelectAIModel, capacityWait=false}) => {
const { t } = useTranslation();
const [sortBy, setSortBy] = React.useState('rerank_score')
const [groupBy, setGroupBy] = React.useState('quality')
@@ -466,7 +436,9 @@ const Candidates = ({rowIndex, rowState, conceptCache, targetCanonical, targetRe
const canFetchMore = hasAnyView
const algoStagesValue = values(rowStage || {}).filter((_, i) => Object.keys(rowStage || {})[i] !== 'recommend')
const areAlgoRun = algoStagesValue.length > 0 && algoStagesValue.every(v => v === 1)
- const { label } = getRowProgressLabel(rowStage, algosSelected);
+ const { label, status: progressStatus } = getRowProgressLabel(rowStage, algosSelected, {t, capacityWait});
+ // Waiting out a busy server, or given up on for now: not a spinner (ocl_issues#2849).
+ const isCapacityStatus = ['capacity_wait', 'throttled'].includes(progressStatus)
const byAlgoScore = sortBy === 'algo_score'
const byRerankScore = sortBy === 'rerank_score'
@@ -628,7 +600,7 @@ const Candidates = ({rowIndex, rowState, conceptCache, targetCanonical, targetRe
{
!areAlgoRun && label &&
- } variant='outlined' color='warning' size='small' label={label} sx={{margin: '0 8px'}} />
+ : } variant='outlined' color='warning' size='small' label={label} sx={{margin: '0 8px'}} />
}
{
// Refresh stays visible even when there are zero candidates so
diff --git a/src/components/map-projects/ConfigurationForm.jsx b/src/components/map-projects/ConfigurationForm.jsx
index dc9782e..bb9630d 100644
--- a/src/components/map-projects/ConfigurationForm.jsx
+++ b/src/components/map-projects/ConfigurationForm.jsx
@@ -61,7 +61,7 @@ const VisuallyHiddenInput = styled('input')({
const deriveCanonicalUrl = relativeUrl => relativeUrl ? `https://ns.openconceptlab.org${relativeUrl}` : ''
-const ConfigurationForm = ({ project, handleFileUpload, file, owner, setOwner, name, setName, description, setDescription, repo, onRepoChange, repoVersion, setRepoVersion, versions, mappedSources, targetSourcesFromRows, algosSelected, setAlgosSelected, sx, algos, validColumns, columns, isValidColumnValue, updateColumn, configure, setConfigure, columnVisibilityModel, setColumnVisibilityModel, onSave, isSaving, candidatesScore, onScoreChange, includeDefaultFilter, setIncludeDefaultFilter, filters, setFilters, locales, isLoadingLocales, setAIAssistantColumns, AIAssistantColumns, inAIAssistantGroup, lookupConfig, setLookupConfig, encoderModel, setEncoderModel, isCoreUser, fullRequestLimits, canSelectAIModel, canSetAIOutputLocale, canBridge, canScispacy, canUseOrgProjects=true, canUseCustomAlgorithms=true, promptTemplates, promptTemplate, onPromptTemplateChange, AIModels, AIModel, setAIModel, namespace, setNamespace, promptOutputLocale, setPromptOutputLocale, inputLocale, setInputLocale, oclLocales, useLexicalVariants, setUseLexicalVariants }) => {
+const ConfigurationForm = ({ project, handleFileUpload, file, owner, setOwner, name, setName, description, setDescription, repo, onRepoChange, repoVersion, setRepoVersion, versions, mappedSources, targetSourcesFromRows, algosSelected, setAlgosSelected, sx, algos, validColumns, columns, isValidColumnValue, updateColumn, configure, setConfigure, columnVisibilityModel, setColumnVisibilityModel, onSave, isSaving, candidatesScore, onScoreChange, includeDefaultFilter, setIncludeDefaultFilter, filters, setFilters, locales, isLoadingLocales, setAIAssistantColumns, AIAssistantColumns, inAIAssistantGroup, lookupConfig, setLookupConfig, encoderModel, setEncoderModel, isCoreUser, requestLimits, canSelectAIModel, canSetAIOutputLocale, canBridge, canScispacy, canUseOrgProjects=true, canUseCustomAlgorithms=true, promptTemplates, promptTemplate, onPromptTemplateChange, AIModels, AIModel, setAIModel, namespace, setNamespace, promptOutputLocale, setPromptOutputLocale, inputLocale, setInputLocale, oclLocales, useLexicalVariants, setUseLexicalVariants }) => {
const { t } = useTranslation();
const user = getCurrentUser()
const isLLMAlgoNotAllowed = !repoVersion?.match_algorithms?.includes('llm')
@@ -438,7 +438,7 @@ const ConfigurationForm = ({ project, handleFileUpload, file, owner, setOwner, n
onChange={setAlgosSelected}
repo={repoVersion}
isCoreUser={isCoreUser}
- fullRequestLimits={fullRequestLimits}
+ requestLimits={requestLimits}
/>
{
isCoreUser &&
diff --git a/src/components/map-projects/MapProject.jsx b/src/components/map-projects/MapProject.jsx
index 173abdf..668a79b 100644
--- a/src/components/map-projects/MapProject.jsx
+++ b/src/components/map-projects/MapProject.jsx
@@ -70,7 +70,8 @@ import pick from 'lodash/pick'
import { OperationsContext } from '../app/LayoutContext';
import APIService, { isTransientNetworkError, retryWithBackoff } from '../../services/APIService';
-import { buildAttributionHeaders, buildConfigSnapshot, summarizeRunCompletion } from '../../services/attribution'
+import { buildAttributionHeaders, buildConfigSnapshot, summarizeRunCompletion, matchesOnOCL } from '../../services/attribution'
+import { CAPACITY_WAIT_CAP_MS, HEAVY_REQUEST_TIMEOUT_MS, LIGHT_REQUEST_TIMEOUT_MS, createCapacityGate, createLimiter, isRetryableError, requestWithCapacityRetry, sleepUnlessCancelled, untilCancelled } from '../../services/capacity'
import { highlightTexts, dropVersion, getCurrentUser, hasAuthGroup, hasCapability, getMapperPreview, getNewProjectBlockReason, downloadObject, currentUserToken, refreshCurrentUserCapabilitiesCache } from '../../common/utils';
import { WHITE, SURFACE_COLORS, TEXT_GRAY } from '../../common/colors';
@@ -86,8 +87,8 @@ import MapProjectDeleteConfirmDialog from './MapProjectDeleteConfirmDialog';
import ConfigurationForm from './ConfigurationForm'
import Controls from './Controls'
import DataGridControls from './DataGridControls'
-import { getPreviewEligibleRowIndexes, getRowsToProcess, spendsMatchQuota, getRowCapByMatchOperations, shouldStopAIStep, getAIRequestIdempotencyKey, getCandidatePoolFingerprint, hasCurrentAnalysis, getScispacyRowResults, getPendingRowLookups, waitForLookups, AI_LOOKUP_WAIT_MS, RERANK_LOOKUP_WAIT_MS } from './autoMatchRows'
-import { createAutosaveScheduler, saveOnLeave, trackSave, whenSaved, installUnloadGuard } from './autosave'
+import { getPreviewEligibleRowIndexes, getRowsToProcess, spendsMatchQuota, getRowCapByMatchOperations, shouldStopAIStep, getAIRequestIdempotencyKey, getCandidatePoolFingerprint, hasCurrentAnalysis, getScispacyRowResults, getPendingRowLookups, waitForLookups, AI_LOOKUP_WAIT_MS, RERANK_LOOKUP_WAIT_MS, RERANK_MAX_IN_FLIGHT, isScispacyWarmingUp } from './autoMatchRows'
+import { createAutosaveScheduler, createLatestSender, saveOnLeave, trackSave, whenSaved, installUnloadGuard } from './autosave'
import MatchSummaryCard from './MatchSummaryCard'
import MappingDecisionResult from './MappingDecisionResult'
import DecisionSelector from './DecisionSelector'
@@ -100,7 +101,7 @@ import Concept from './Concept'
import ImportToCollection from './ImportToCollection'
import ProjectLogs from './ProjectLogs';
import { useAlgos, ensureConceptIdentity } from './algorithms'
-import { applyRequestSettings, getRequestSettings, hasFullRequestLimits } from './requestLimits'
+import { applyRequestSettings, getRequestLimits, getRequestSettings } from './requestLimits'
import { countQualityBuckets, createBestCandidateScoreCache, filterRowsByQualityBucket, getRowQualities } from './rowQuality'
import AutoMatchDialog from './AutoMatchDialog'
import PreviewLimitDialog from './PreviewLimitDialog'
@@ -108,7 +109,7 @@ import QuotaDialog from '../common/QuotaDialog'
import { getQuotaError } from '../common/quotaErrors'
import MapperQuotaChip from './MapperQuotaChip'
import { getPreviewLimitError, isAIPreviewLimitError } from './previewLimits'
-import { formatRowNumbers, runMatchBatch, requestSingleMatch } from './matchBatch'
+import { formatRowNumbers, runMatchBatch, requestSingleMatch, runWithConcurrency, INTERACTIVE_WAIT_CAP_MS } from './matchBatch'
import { getAIAssistantChoices, getModelUsed, getPromptTemplateKey, createLatestRequestGate } from './aiVisibility'
import { DEFAULT_ENCODER_MODEL } from './rerankerModels'
import { normalizeAlgorithmInvocation, hasSuccessfulAlgorithmResponse, getAlgorithmStagesFromResponses, lookupStatusRank, buildRecommendableConceptEntry, stripConstantClassAndDatatype, buildLookupConceptUrl, serializeCandidates, toStoredRowMatchState } from './normalizers'
@@ -226,6 +227,28 @@ const MapProject = () => {
const autoMatchRunGateRef = React.useRef(null)
if(!autoMatchRunGateRef.current)
autoMatchRunGateRef.current = createLatestRequestGate()
+ // The latest run's ticket, for waits that must end when that run stops.
+ const autoMatchRunTicketRef = React.useRef(null)
+ // Waiting out a busy server (ocl_issues#2849). One gate per server: a 429 on
+ // any request holds back this tab's others to that server, and the Auto
+ // Match scheduler. At most RERANK_MAX_IN_FLIGHT $rerank calls at once, the
+ // per-row reranks during a run included.
+ const capacityGatesRef = React.useRef(new Map())
+ const rerankLimiterRef = React.useRef(null)
+ if(!rerankLimiterRef.current)
+ rerankLimiterRef.current = createLimiter(RERANK_MAX_IN_FLIGHT)
+ // 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)
+ // 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({})
+ // One logs POST at a time per project URL (see createLatestSender).
+ const logsSendersRef = React.useRef(new Map())
+ // The row a bridge $match is for, set around the call into the Bridge Match
+ // component (see getBridgeMatchService).
+ const bridgeCallContextRef = React.useRef(null)
// ai_assistant.calls is a one-time allowance (no reset, R2) - once exhausted it
// stays exhausted for the rest of the session, so this is never reset per-run.
const aiQuotaExhaustedRef = React.useRef(false);
@@ -566,8 +589,9 @@ const MapProject = () => {
// in-flight check"). Replaces the legacy "wait for every algo to
// complete" trigger with: any ConceptRow with rerank_score === undefined
// makes its row eligible; debounce coalesces rapid algo completions;
- // in-flight set prevents double-firing.
- const inFlightRerankRef = React.useRef(new Set())
+ // in-flight map (row -> a promise that settles when its rerank ends)
+ // prevents double-firing.
+ const inFlightRerankRef = React.useRef(new Map())
const rerankDebounceRef = React.useRef({})
const rerankRerunNeededRef = React.useRef(new Set())
// Forward-ref pointer for mergeIntoRowMatchState (declared earlier in
@@ -585,8 +609,9 @@ const MapProject = () => {
const OCL_ONLINE_API_URL = window.OCL_ONLINE_API_URL || process.env.OCL_ONLINE_API_URL
const inAIAssistantGroup = Boolean(hasCapability(user, 'users.mapper_ai_assistant') && AI_ASSISTANT_API_URL)
const isCoreUser = hasAuthGroup(user, 'core_user')
- // Batch size and concurrent requests per algorithm (ocl_online#274)
- const fullRequestLimits = hasFullRequestLimits(user)
+ // Batch size and concurrent requests per algorithm, by the user's tier
+ // (ocl_online#274, ocl_issues#2849)
+ const requestLimits = getRequestLimits(user)
const mapperPreview = getMapperPreview()
// Only staff choose the AI model and prompt template (canSelectAIModel gates
// both pickers). Core, early-access and unlimited-AI users keep the output
@@ -1637,9 +1662,17 @@ const MapProject = () => {
user: user.username || user.id,
extras: isAutoSave ? {reasons: options.reasons || []} : (isUpdate ? undefined : {project: response.data})
}
- const savedProjectLogs = [saveLog, ...projectLogsForSave]
- projectLogsRef.current = savedProjectLogs
- setProjectLogs(savedProjectLogs)
+ // Entries logged while the save was in flight are kept. A save that
+ // lands after the user opened another project uses the logs as they
+ // were when it started, and leaves the open project's alone
+ // (ocl_issues#2849).
+ const isOtherProjectOpen = isUpdate && String(projectIdRef.current) !== String(response.data.id)
+ const savedProjectLogs = [saveLog, ...(isOtherProjectOpen ? projectLogsForSave : projectLogsRef.current)]
+ const savedRowLogs = isOtherProjectOpen ? rowLogsForSave : logsRef.current
+ if(!isOtherProjectOpen) {
+ projectLogsRef.current = savedProjectLogs
+ setProjectLogs(savedProjectLogs)
+ }
if(options.closeConfigure !== false) {
configSnapshotOnOpenRef.current = null
setConfigure(false)
@@ -1654,7 +1687,7 @@ const MapProject = () => {
if(!isAutoSave)
baseSetAlert({severity: 'success', message: t('map_project.successfully_saved'), duration: 2000})
- APIService.new().overrideURL(response.data.url).appendToUrl('logs/').post({logs: {row_logs: rowLogsForSave, project_logs: savedProjectLogs}}).then(() => {})
+ postProjectLogs(response.data.url, {row_logs: savedRowLogs, project_logs: savedProjectLogs})
} else if(options.leaving && status !== 401) {
// The user has left the project, so no dialog can show: say so once.
baseSetAlert({severity: 'error', message: t('map_project.leave_save_failed', {name: name || project?.name, detail: errorData?.detail || t('unknown_error')}), duration: 12000})
@@ -1690,11 +1723,112 @@ const MapProject = () => {
projectLogsRef.current = newLogs
setProjectLogs(newLogs)
if(project?.url)
- APIService.new().overrideURL(project.url).appendToUrl('logs/').post({logs: {row_logs: logsRef.current, project_logs: newLogs}}).then(() => {})
+ postProjectLogs(project.url, {row_logs: logsRef.current, project_logs: newLogs})
+ }
+
+ // POSTs the project's logs, waiting out a busy server. One POST at a time
+ // per project: each sends the whole log, so an older one that waited out a
+ // 429 mustn't land after a newer one, and a retry sends the latest log
+ // triggered for that project, never another project's (ocl_issues#2849).
+ // Best effort, like before.
+ const postProjectLogs = (url, logs) => {
+ if(!url)
+ return
+ if(!logsSendersRef.current.has(url))
+ logsSendersRef.current.set(url, createLatestSender(getPayload => requestWithCapacityRetry(
+ () => APIService.new().overrideURL(url).appendToUrl('logs/').request('POST', getPayload(), null, {handlesThrottle: true, timeout: LIGHT_REQUEST_TIMEOUT_MS}),
+ {maxWaitMs: INTERACTIVE_WAIT_CAP_MS}
+ )))
+ logsSendersRef.current.get(url).trigger({logs})
}
const scheduleAutoSave = reason => autosaveSchedulerRef.current.schedule(reason)
+ // ── Waiting out a busy server (ocl_issues#2849) ───────────────────────────
+ // The gate for a request's server, shared by every request this tab sends it.
+ const getCapacityGate = url => {
+ let origin
+ try {
+ origin = new URL(url, window.location.href).origin
+ } catch (_) {
+ origin = String(url)
+ }
+ if(!capacityGatesRef.current.has(origin))
+ capacityGatesRef.current.set(origin, createCapacityGate())
+ return capacityGatesRef.current.get(origin)
+ }
+
+ const syncCapacityWaits = () => {
+ const rows = {}
+ capacityWaitsRef.current.forEach(rowIndexes => rowIndexes.forEach(index => { rows[index] = true }))
+ setCapacityWaitRows(capacityWaitsRef.current.size ? rows : null)
+ }
+
+ // 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
+ // each row's log with the server's capacity headers.
+ const trackCapacityWait = (rowIndexes, logExtras = {}) => {
+ const waitId = {}
+ let logged = false
+ return {
+ onWait: ({reason, retryAfterMs, capacity}) => {
+ if(reason === 'error')
+ return
+ capacityWaitsRef.current.set(waitId, rowIndexes)
+ 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} : {})}
+ }, index))
+ },
+ onWaitEnd: () => {
+ if(capacityWaitsRef.current.delete(waitId))
+ syncCapacityWaits()
+ },
+ }
+ }
+
+ // The server's X-OCL-Capacity-* headers are only logged for now; adapting to
+ // them is OpenConceptLab/ocl_online#340.
+ const noteCapacity = capacity => console.debug('OCL capacity', capacity)
+
+ // Stop for the run in progress: true once the user stops it, or a newer run
+ // starts. Waits in a run poll it.
+ const getRunStopCheck = () => {
+ const ticket = autoMatchRunTicketRef.current
+ return () => Boolean(abortRef.current || ticket !== autoMatchRunTicketRef.current)
+ }
+
+ // A result that lands after the user stopped its run is kept (a graceful
+ // stop), so save it too: the stopped run's own autosave may have gone
+ // already (ocl_issues#2849).
+ const saveLateRunResult = () => {
+ if(abortRef.current)
+ scheduleAutoSave('auto_match_stopped')
+ }
+
+ const noteThrottledRunRow = (algoId, index) => {
+ throttledRunRowsRef.current = {
+ ...throttledRunRowsRef.current,
+ [algoId]: uniq([...(throttledRunRowsRef.current[algoId] || []), index]),
+ }
+ }
+
+ // Once a request for an algorithm (or rerank) stayed refused for the whole
+ // cap, the run stops asking for it: the rest of its rows end throttled at
+ // once, "not run, retry", instead of each waiting out the cap again.
+ const isRunThrottled = algoId => Boolean(throttledRunRowsRef.current[algoId]?.length)
+ const markRowsThrottled = (indexes, algoId, logExtras = {}) => indexes.forEach(index => {
+ markAlgo(index, algoId, -4)
+ log({action: algoId === 'rerank' ? 'rerank_throttled' : 'algo_throttled', description: t('map_project.row_throttled'), extras: {...logExtras, skipped: true}}, index)
+ noteThrottledRunRow(algoId, index)
+ })
+
// ── ocl_online#105 Phase 5: AutomatchRun attribution ──────────────────────
// Build the X-OCL-Request-Source + X-OCL-Event-Metadata headers for a single
// per-row backend call. Attribution is by explicit ORIGIN, not ambient: only
@@ -1702,11 +1836,13 @@ const MapProject = () => {
// call (panel open, single rerun, per-row AI) that overlaps a running
// auto-match is correctly stamped 'mapper-ui-manual' — the active-run ref
// alone can't distinguish concurrent manual traffic from run traffic.
- const attrHeaders = ({rowIndex, rowIndices, batchSize, algorithmId, clientAttemptN, isRunTraffic = false} = {}) =>
+ // capacityAware: the call waits out a capacity 429 (OpenConceptLab/ocl_online#275). Set it on OCL's gated calls
+ // (semantic or reranked $match, $rerank), never on a custom algorithm's own server.
+ const attrHeaders = ({rowIndex, rowIndices, batchSize, algorithmId, clientAttemptN, isRunTraffic = false, capacityAware = false} = {}) =>
buildAttributionHeaders({
runId: isRunTraffic ? (automatchRunRef.current?.id ?? null) : null,
projectId: params.projectId,
- rowIndex, rowIndices, batchSize, algorithmId, clientAttemptN,
+ rowIndex, rowIndices, batchSize, algorithmId, clientAttemptN, capacityAware,
})
// Create the AutomatchRun system-of-record at run start (oclapi2#876), which
@@ -1737,14 +1873,38 @@ const MapProject = () => {
}),
}
try {
- const response = await APIService.new().overrideURL(project.url).appendToUrl('auto-match-runs/')
- .post(body, null, {}, undefined, true)
- const status = response?.response?.status || response?.status
- const data = response?.response?.data || response?.data
+ // A busy server's 429 is waited out, and Stop ends the wait, and the
+ // wait for an answer. Nothing else is retried: after a 502 or a timeout
+ // the run may exist already, and a retry would open a second one
+ // (ocl_issues#2849).
+ const isStopped = getRunStopCheck()
+ const request = requestWithCapacityRetry(
+ () => APIService.new().overrideURL(project.url).appendToUrl('auto-match-runs/').request('POST', body, null, {handlesThrottle: true, timeout: LIGHT_REQUEST_TIMEOUT_MS}),
+ {maxWaitMs: INTERACTIVE_WAIT_CAP_MS, isCancelled: isStopped, isRetryable: () => false, onCapacity: noteCapacity}
+ )
+ const outcome = await untilCancelled(request, isStopped)
+ if(!outcome.settled) {
+ // Stopped before the server answered: close a run it opens anyway.
+ // Idempotent, so retried like the normal close; never rejects.
+ request.then(late => {
+ const lateId = late.ok && late.response?.data?.id
+ if(lateId)
+ requestWithCapacityRetry(
+ () => APIService.new().overrideURL('/auto-match-runs/' + lateId + '/')
+ .request('PATCH', {completed_rows: 0, failed_rows: 0, completion_status: 'cancelled'}, null, {handlesThrottle: true, timeout: LIGHT_REQUEST_TIMEOUT_MS}),
+ {maxWaitMs: INTERACTIVE_WAIT_CAP_MS}
+ )
+ })
+ return
+ }
+ const result = outcome.value
+ const response = result.ok ? result.response : result.error?.response
+ const status = response?.status
+ const data = response?.data
const id = data?.id
if(id) {
automatchRunRef.current = {id, algoIds: map(selectedAlgos, 'id')}
- return
+ return automatchRunRef.current
}
if(status === 403 && data?.error_code) {
abortRef.current = true
@@ -1753,32 +1913,37 @@ const MapProject = () => {
setPreviewLimit({errorCode: data.error_code, limit: data.limit, used: data.used})
return
}
- projectLog({action: 'automatch_run_create_failed', extras: {status: status || 'no-id'}})
+ projectLog({action: 'automatch_run_create_failed', extras: {status: status || 'no-id', ...(result.ok ? {} : {reason: result.reason, error: result.error?.message || null})}})
} catch (err) {
projectLog({action: 'automatch_run_create_failed', extras: {error: err?.message || 'unknown'}})
}
}
// PATCH the run to a terminal status at run end. completed/failed are derived
- // from the per-row algo stages (rowStageRef: -2 failed, 1 done, -1 not run)
+ // from the per-row algo stages (rowStageRef: -4 throttled, -2 failed, 1 done, -1 not run)
// over the rows we set out to process; a row counts as failed only when EVERY
// attempted algo failed for it. Never throws — a failed PATCH (or a run that
// was never created) must not surface to the user.
- const completeAutomatchRun = async (selectedAlgos, intendedRows) => {
- const run = automatchRunRef.current
- automatchRunRef.current = null
+ // Each run closes its own record, passed in: by the time a stopped run gets
+ // here, a newer run's record may be the current one (ocl_issues#2849).
+ const completeAutomatchRun = async (selectedAlgos, intendedRows, run, aborted) => {
+ if(automatchRunRef.current === run)
+ automatchRunRef.current = null
if(!run?.id) return
const payload = summarizeRunCompletion({
rowStages: rowStageRef.current || [],
rowIndices: map(intendedRows, '__index'),
algoIds: run.algoIds?.length ? run.algoIds : map(selectedAlgos, 'id'),
- aborted: abortRef.current,
+ aborted,
stoppedForQuota: Boolean(matchQuotaStopRef.current),
})
- try {
- await APIService.new().overrideURL('/auto-match-runs/' + run.id + '/')
- .request('PATCH', payload)
- } catch (_) { /* attribution is best-effort; never block on the PATCH */ }
+ // Closing the run is idempotent, so a busy server or a gateway error is
+ // retried; Stop doesn't end it, so a stopped run is closed too
+ // (ocl_issues#2849). requestWithCapacityRetry never rejects.
+ await requestWithCapacityRetry(
+ () => APIService.new().overrideURL('/auto-match-runs/' + run.id + '/').request('PATCH', payload, null, {handlesThrottle: true, timeout: LIGHT_REQUEST_TIMEOUT_MS}),
+ {maxWaitMs: INTERACTIVE_WAIT_CAP_MS}
+ )
}
const fetchRepo = (url, _repo) => APIService.new().overrideURL(url).get().then(response => setRepo(response.data?.id ? response.data : _repo))
@@ -1923,6 +2088,42 @@ const MapProject = () => {
return service
}
+ // The service handed to the Bridge Match component (react-bridge-match),
+ // which sends the bridge $match through its post(). APIService.post never
+ // settled on a 429, which froze a run's bridge step, and resolved a 5xx as
+ // its body. This post() waits out 429s, retries transient failures, and
+ // always settles with what the component reads: the response, or a body
+ // with a detail (ocl_issues#2849). A new instance each render, as before:
+ // the component points it at its own URL for a GET.
+ const getBridgeMatchService = () => {
+ const service = getMatchAPIService()
+ service.post = (body, token, headers, query) => {
+ const context = bridgeCallContextRef.current || {}
+ context.posted = true
+ return requestWithCapacityRetry(() => service.request('POST', body, token, {headers: headers || {}, query, handlesThrottle: true, timeout: HEAVY_REQUEST_TIMEOUT_MS}), {
+ gate: getCapacityGate(service.URL),
+ isCancelled: context.isCancelled || (() => false),
+ maxWaitMs: context.isBulk ? CAPACITY_WAIT_CAP_MS : INTERACTIVE_WAIT_CAP_MS,
+ onCapacity: noteCapacity,
+ ...(isNumber(context.rowIndex) ? trackCapacityWait([context.rowIndex], getAlgoLogExtras(bridgeAlgo)) : {}),
+ }).then(result => {
+ if(result.ok)
+ return result.response
+ if(result.reason === 'throttled')
+ return {detail: t('map_project.algorithm_throttled'), status: 429, throttled: true}
+ if(result.reason === 'cancelled')
+ return {detail: 'cancelled', cancelled: true}
+ const data = result.error?.response?.data
+ const hasServerBody = data && typeof data === 'object' && !Array.isArray(data) && (data.detail || data.error_code)
+ return hasServerBody ? data : {
+ detail: t('map_project.match_request_failed', {error: result.error?.response?.data?.detail || result.error?.message || t('unknown_error')}),
+ status: result.error?.response?.status || null,
+ }
+ })
+ }
+ return service
+ }
+
// Pick the highest-scoring target-repo ConceptRow for a given rowIndex
// from the unified-model state. Returns a RowView ({candidate,
// conceptDefinition, conceptRow, bridgeConceptDefinition?}) or null.
@@ -2057,12 +2258,19 @@ const MapProject = () => {
// reported at the end of the run so the user knows which rows to re-run.
const failedMatchRows = {}
const runTicket = autoMatchRunGateRef.current.next()
+ autoMatchRunTicketRef.current = runTicket
const isCurrentRun = () => autoMatchRunGateRef.current.isCurrent(runTicket)
+ // This run's Stop: the user stopped it, or a newer run started (which
+ // resets abortRef). Every phase checks this, not abortRef alone, so a
+ // stopped run can't carry on once the next one starts (ocl_issues#2849).
+ const isRunStopped = () => Boolean(abortRef.current || !isCurrentRun())
+ throttledRunRowsRef.current = {}
// Function to process a single batch
const processBatch = async (_repo, rowBatch, algo) => {
- if (abortRef.current) {
- setLoadingMatches(false)
+ if (isRunStopped()) {
+ if(isCurrentRun())
+ setLoadingMatches(false)
return []
};
if (matchQuotaStopRef.current && spendsMatchQuota(algo, {canBridge}))
@@ -2091,7 +2299,8 @@ const MapProject = () => {
const logExtras = getAlgoLogExtras(algo)
// service.request rejects on a network error or a non-2xx; service.post
// would resolve, and the failed batch would pass as a success with no
- // candidates (ocl_online#257). runMatchBatch retries transient failures.
+ // candidates (ocl_online#257). runMatchBatch retries transient failures,
+ // and waits out a busy server's 429s (ocl_issues#2849).
return runMatchBatch({
rowIndexes,
send: attempt => service.request(
@@ -2099,14 +2308,18 @@ const MapProject = () => {
payload,
(algo.type === 'custom' && algo.url && algo.token) ? algo.token : null,
{
- headers: attrHeaders({rowIndices: rowIndexes, batchSize: rowBatch.length, algorithmId: algo.id, clientAttemptN: attempt + 1, isRunTraffic: true}),
+ headers: attrHeaders({rowIndices: rowIndexes, batchSize: rowBatch.length, algorithmId: algo.id, clientAttemptN: attempt + 1, isRunTraffic: true, capacityAware: matchesOnOCL(algo)}),
query: {
includeSearchMeta: true,
...(algo.query_params || {}),
...extraParams
- }
+ },
+ handlesThrottle: true,
+ timeout: HEAVY_REQUEST_TIMEOUT_MS,
}
),
+ retryOptions: {gate: getCapacityGate(service.URL), onCapacity: noteCapacity},
+ ...trackCapacityWait(rowIndexes, logExtras),
// A batch that outlives its run must not touch the next run's stages.
setStage: (index, stage) => {
if(isCurrentRun())
@@ -2118,9 +2331,19 @@ const MapProject = () => {
if(!previewLimit)
failedMatchRows[algo.id] = [...(failedMatchRows[algo.id] || []), index]
},
+ // Still refused after the long cap: not run, and not a failure.
+ onRowThrottled: (index, {attempts, waitedMs}) => {
+ log({action: 'algo_throttled', description: t('map_project.row_throttled'), extras: {...logExtras, attempts, waited_ms: waitedMs}}, index)
+ if(isCurrentRun())
+ noteThrottledRunRow(algo.id, index)
+ },
// Sets matchQuotaStopRef, which stops the rest of the $match requests
// but lets the run finish with what it has.
- onPreviewLimit: err => handlePreviewLimitError(err),
+ // Not once a newer run started: its quota state is its own.
+ onPreviewLimit: err => {
+ if(isCurrentRun())
+ handlePreviewLimitError(err)
+ },
// A batch waiting to retry stops on a cancel, once another batch has
// hit the match quota, or once a newer run has started.
shouldStop: () => Boolean(
@@ -2132,100 +2355,82 @@ const MapProject = () => {
// Function to handle concurrency
const processWithConcurrency = async (_repo, algo, _rows) => {
// A project's saved settings run within this user's limits.
- const { batchSize, concurrentRequests } = getRequestSettings(algo, fullRequestLimits)
- const CHUNK_SIZE = batchSize // Number of rows per batch
- const MAX_CONCURRENT_REQUESTS = concurrentRequests; // Number of parallel API requests allowed
- const rowChunks = chunk(_rows, CHUNK_SIZE);
-
- const queue = rowChunks.slice(); // Copy of all chunks to be processed
- const activeRequests = new Set();
-
- while (queue.length > 0 || activeRequests.size > 0) {
- // Fill activeRequests up to MAX_CONCURRENT_REQUESTS
- while (queue.length > 0 && activeRequests.size < MAX_CONCURRENT_REQUESTS) {
- if (abortRef.current) {
- setLoadingMatches(false)
+ const { batchSize, concurrentRequests } = getRequestSettings(algo, requestLimits)
+ // While a 429 holds the server's gate, no new batch goes out: the queue
+ // waits instead of draining into refusals (ocl_issues#2849).
+ const isQuotaStopped = () => Boolean(matchQuotaStopRef.current && spendsMatchQuota(algo, {canBridge}))
+ const finished = await runWithConcurrency(chunk(_rows, batchSize), {
+ concurrency: concurrentRequests,
+ gate: getCapacityGate(getMatchAPIService(algo).URL),
+ maxHoldMs: CAPACITY_WAIT_CAP_MS,
+ shouldAbort: isRunStopped,
+ shouldSkipRest: () => isQuotaStopped() || isRunThrottled(algo.id),
+ // A quota stop leaves the rest not run, as before; a server that
+ // stayed too busy leaves them throttled.
+ onSkipped: rowBatches => {
+ if(!isQuotaStopped())
+ markRowsThrottled(map(flatten(rowBatches), '__index'), algo.id, getAlgoLogExtras(algo))
+ },
+ run: rowBatch => processBatch(_repo, rowBatch, algo).then((data) => {
+ // A batch that answers after a newer run started belongs to no run.
+ if(!isCurrentRun())
return
- };
- if (matchQuotaStopRef.current && spendsMatchQuota(algo, {canBridge})) {
- queue.length = 0
- break
- }
- const rowBatch = queue.shift();
- const promise = processBatch(_repo, rowBatch, algo).then((data) => {
- // Populate rowMatchState before any consumer (setStateViews /
- // setAutoMatched) tries to read it. The per-row fetch flow
- // routes through `onResponse` which already calls
- // mergeIntoRowMatchState; bulk processBatch doesn't, so we
- // run the normalizer here.
- if(Array.isArray(data) && data.length) {
- const projectCtx = buildProjectContext()
- if(projectCtx) {
- const algoCfg = ensureConceptIdentity(algo)
- if(algoCfg) {
- data.forEach(rowEntry => {
- const idx = rowEntry?.row?.__index
- if(!isNumber(idx)) return
- mergeIntoRowMatchState(idx, normalizeAlgorithmInvocation(rowEntry, {
- algorithmId: algo.id,
- algorithmConfig: algoCfg,
- projectContext: projectCtx,
- rowIndex: idx,
- // Bulk auto-match's $match request sends reranker=!isMultiAlgo
- // (line 1433). Trust the server's normalized score only when
- // that flag was true — otherwise it's a per-algo native score
- // (e.g. FAISS similarity × 100) masquerading as a rerank.
- trustServerRerank: !isMultiAlgo
- }), {isRunTraffic: true})
- })
- }
+ // Populate rowMatchState before any consumer (setStateViews /
+ // setAutoMatched) tries to read it. The per-row fetch flow
+ // routes through `onResponse` which already calls
+ // mergeIntoRowMatchState; bulk processBatch doesn't, so we
+ // run the normalizer here.
+ if(Array.isArray(data) && data.length) {
+ const projectCtx = buildProjectContext()
+ if(projectCtx) {
+ const algoCfg = ensureConceptIdentity(algo)
+ if(algoCfg) {
+ data.forEach(rowEntry => {
+ const idx = rowEntry?.row?.__index
+ if(!isNumber(idx)) return
+ mergeIntoRowMatchState(idx, normalizeAlgorithmInvocation(rowEntry, {
+ algorithmId: algo.id,
+ algorithmConfig: algoCfg,
+ projectContext: projectCtx,
+ rowIndex: idx,
+ // Bulk auto-match's $match request sends reranker=!isMultiAlgo
+ // (line 1433). Trust the server's normalized score only when
+ // that flag was true — otherwise it's a per-algo native score
+ // (e.g. FAISS similarity × 100) masquerading as a rerank.
+ trustServerRerank: !isMultiAlgo
+ }), {isRunTraffic: true})
+ })
}
}
- if(!isMultiAlgo)
- setStateViews(data, _repo)
- setMatchedConcepts(prev => [...prev, ...data]);
- activeRequests.delete(promise); // Remove from active set after completion
- });
- activeRequests.add(promise);
- }
-
- // Wait for at least one request to complete before continuing
- if (activeRequests.size > 0)
- await Promise.race(activeRequests);
- }
+ }
+ if(!isMultiAlgo)
+ setStateViews(data, _repo)
+ setMatchedConcepts(prev => [...prev, ...data]);
+ if(Array.isArray(data) && data.length)
+ saveLateRunResult()
+ }),
+ })
+ if(!finished && isCurrentRun())
+ setLoadingMatches(false)
};
+ // Reranks the run's rows, at most maxConcurrent at once, through the same
+ // limiter as every other rerank. Stop returns at once.
const processRerankWithConcurrency = async (_rows, maxConcurrent = 2) => {
- const queue = _rows.slice();
- const activeRequests = new Set();
-
- while (queue.length > 0 || activeRequests.size > 0) {
- while (queue.length > 0 && activeRequests.size < maxConcurrent) {
- if (abortRef.current) {
- setLoadingMatches(false)
- return
- }
- if(rerankQuotaStopRef.current) {
- queue.length = 0
- break
- }
-
- const row = queue.shift();
- const promise = rerank(row.__index, true)
- .catch(() => null)
- .finally(() => activeRequests.delete(promise));
- activeRequests.add(promise);
- }
-
- if(activeRequests.size > 0)
- await Promise.race(activeRequests);
- }
+ const finished = await runWithConcurrency(_rows, {
+ concurrency: maxConcurrent,
+ shouldAbort: isRunStopped,
+ shouldSkipRest: () => Boolean(rerankQuotaStopRef.current),
+ run: row => rerank(row.__index, true).catch(() => null),
+ })
+ if(!finished && isCurrentRun())
+ setLoadingMatches(false)
};
// Capped users' algorithms carry the settings they run with, so the
// AutomatchRun's config snapshot records what ran (ocl_online#274).
let _selectedAlgos = map(
filter(algosSelected, algo => selectedAlgos.includes(algo.id)),
- algo => applyRequestSettings(algo, fullRequestLimits)
+ algo => applyRequestSettings(algo, requestLimits)
)
let subActions = [...map(_selectedAlgos, algo => algo.name || algo.id)]
subActions.push('reranker')
@@ -2282,7 +2487,10 @@ const MapProject = () => {
// ocl_online#105 Phase 5: open the run record, then guarantee it is
// closed out (completed / partial / failed / cancelled) via the finally,
// even on cancel or an unexpected throw mid-pipeline.
- await createAutomatchRun(_selectedAlgos, rowsToProcess)
+ const runRecord = await createAutomatchRun(_selectedAlgos, rowsToProcess)
+ // False when the run was refused (a preview limit) or stopped before it
+ // began: it changed nothing to save.
+ const runStarted = !isRunStopped()
try {
bulkMatchAlgoIdsRef.current = map(_selectedAlgos, 'id')
// Reset all algo stages to -1 for every row before starting so that
@@ -2294,6 +2502,9 @@ const MapProject = () => {
const rowId = row.__index
const rowState = { ...(next[rowId] || {}) }
_selectedAlgos.forEach(algo => { rowState[algo.id] = -1 })
+ // A rerank an earlier run left throttled is this run's to retry.
+ if(rowState.rerank === -4)
+ rowState.rerank = -1
next[rowId] = rowState
})
rowStageRef.current = next
@@ -2302,14 +2513,14 @@ const MapProject = () => {
isBulkMatchRunningRef.current = true
try {
for(const algo of _selectedAlgos) {
- if(abortRef.current) break
+ if(isRunStopped()) break
if(matchQuotaStopRef.current && spendsMatchQuota(algo, {canBridge})) continue
if(['custom', 'ocl-search', 'ocl-semantic'].includes(algo.type))
await processWithConcurrency(repo, algo, rowsToProcess)
else if(['ocl-bridge', 'ocl-ciel-bridge'].includes(algo.type) && canBridge)
- await fetchBulkBridgeCandidates(rowsToProcess, algo)
+ await fetchBulkBridgeCandidates(rowsToProcess, algo, isRunStopped)
else if(algo.type === 'ocl-scispacy' && canScispacy)
- await fetchBulkScispacyCandidates(rowsToProcess, algo)
+ await fetchBulkScispacyCandidates(rowsToProcess, algo, isRunStopped)
}
} finally {
isBulkMatchRunningRef.current = false
@@ -2322,37 +2533,57 @@ const MapProject = () => {
rowsToProcess
if(_selectedAlgos.length)
await processRerankWithConcurrency(finishedRows, 2)
- if(inAIAssistantGroup && autoRunAIAnalysis) {
+ const runAIStep = Boolean(inAIAssistantGroup && autoRunAIAnalysis && !isRunStopped())
+ if(runAIStep)
await new Promise(resolve => setTimeout(resolve, 1000))
- await runBulkAIAnalysis(finishedRows)
- } else {
+ if(runAIStep && !isRunStopped()) {
+ await runBulkAIAnalysis(finishedRows, {isStopped: isRunStopped, isCurrent: isCurrentRun})
+ } else if(isCurrentRun()) {
+ // A newer run's loading state is its own.
setIsLoadingInDecisionView(false)
setLoadingMatches(false)
setEndMatchingAt(moment())
}
- const failedAlgoIds = keys(failedMatchRows)
- if(!abortRef.current && failedAlgoIds.length) {
+ const getAlgoLabel = algoId => {
+ if(algoId === 'rerank')
+ return t('map_project.rerank_step')
+ const name = find(_selectedAlgos, {id: algoId})?.name
+ return isString(name) && name ? name : algoId
+ }
+ const getRowsByAlgo = rowsPerAlgo => {
const rowsByAlgo = {}
- failedAlgoIds.forEach(algoId => { rowsByAlgo[algoId] = uniq(failedMatchRows[algoId]).sort((a, b) => a - b) })
- const getAlgoLabel = algoId => {
- const name = find(_selectedAlgos, {id: algoId})?.name
- return isString(name) && name ? name : algoId
- }
+ keys(rowsPerAlgo).forEach(algoId => { rowsByAlgo[algoId] = uniq(rowsPerAlgo[algoId]).sort((a, b) => a - b) })
+ return rowsByAlgo
+ }
+ const describeRows = rowsByAlgo => keys(rowsByAlgo).map(algoId => t('map_project.auto_match_rows_failed_algorithm', {
+ algorithm: getAlgoLabel(algoId),
+ rows: formatRowNumbers(rowsByAlgo[algoId]),
+ })).join('; ')
+ const endOfRunNotices = []
+ const failedAlgoIds = keys(failedMatchRows)
+ if(!isRunStopped() && failedAlgoIds.length) {
+ const rowsByAlgo = getRowsByAlgo(failedMatchRows)
projectLog({action: 'auto_match_rows_failed', extras: {row_indexes_by_algorithm: rowsByAlgo}})
- setAlert({
- severity: 'warning',
- message: t('map_project.auto_match_rows_failed', {
- count: uniq(flatten(values(rowsByAlgo))).length,
- details: failedAlgoIds.map(algoId => t('map_project.auto_match_rows_failed_algorithm', {
- algorithm: getAlgoLabel(algoId),
- rows: formatRowNumbers(rowsByAlgo[algoId]),
- })).join('; '),
- }),
- })
+ endOfRunNotices.push(t('map_project.auto_match_rows_failed', {
+ count: uniq(flatten(values(rowsByAlgo))).length,
+ details: describeRows(rowsByAlgo),
+ }))
+ }
+ // Rows the server stayed too busy for weren't run: say so apart from
+ // the failures (ocl_issues#2849).
+ if(!isRunStopped() && keys(throttledRunRowsRef.current).length) {
+ const rowsByAlgo = getRowsByAlgo(throttledRunRowsRef.current)
+ projectLog({action: 'auto_match_rows_throttled', extras: {row_indexes_by_algorithm: rowsByAlgo}})
+ endOfRunNotices.push(t('map_project.auto_match_rows_throttled', {
+ count: uniq(flatten(values(rowsByAlgo))).length,
+ details: describeRows(rowsByAlgo),
+ }))
}
- if(!abortRef.current && matchQuotaStopRef.current)
+ if(endOfRunNotices.length)
+ setAlert({severity: 'warning', message: endOfRunNotices.join(' ')})
+ if(!isRunStopped() && matchQuotaStopRef.current)
projectLog({action: 'auto_match_stopped_for_quota', extras: {reason: matchQuotaStopRef.current}})
- if(!abortRef.current)
+ if(!isRunStopped())
projectLog({
action: 'auto_match_finished',
extras: {
@@ -2367,10 +2598,12 @@ const MapProject = () => {
} : {})
}
})
- if(!abortRef.current)
- scheduleAutoSave('auto_match')
} finally {
- await completeAutomatchRun(_selectedAlgos, rowsToProcess)
+ // Save what the run matched, also when it was stopped or failed part
+ // way (ocl_issues#2849).
+ if(runStarted)
+ scheduleAutoSave(isRunStopped() ? 'auto_match_stopped' : 'auto_match')
+ await completeAutomatchRun(_selectedAlgos, rowsToProcess, runRecord, isRunStopped())
refreshMapperQuotaCache()
}
}, 1000)
@@ -2380,7 +2613,9 @@ const MapProject = () => {
conceptCacheRef.current = conceptCache;
}, [conceptCache]);
- const runBulkAIAnalysis = async (_rows) => {
+ // isStopped: the run's Stop; isCurrent: false once a newer run started, whose
+ // loading state is its own (ocl_issues#2849).
+ const runBulkAIAnalysis = async (_rows, {isStopped = () => abortRef.current, isCurrent = () => true} = {}) => {
GAService.recordActionEvent('MapProject', 'ai_assistant_run', undefined, {
mode: 'bulk',
row_count: _rows.length
@@ -2389,7 +2624,20 @@ const MapProject = () => {
setBulkAIAnalysisStartedAt(moment())
let resolvedPromptTemplate
try {
- resolvedPromptTemplate = await resolvePromptTemplateForInvocation()
+ // Stop doesn't wait on a call that never answers (ocl_issues#2849); the
+ // AI Assistant calls' own retries are OpenConceptLab/ocl_online#206.
+ const outcome = await untilCancelled(resolvePromptTemplateForInvocation(), isStopped)
+ if(!outcome.settled) {
+ if(isCurrent()) {
+ const now = moment()
+ setBulkAIAnalysisEndedAt(now)
+ setEndMatchingAt(now)
+ setLoadingMatches(false)
+ setIsLoadingInDecisionView(false)
+ }
+ return
+ }
+ resolvedPromptTemplate = outcome.value
} catch (err) {
const now = moment()
setBulkAIAnalysisEndedAt(now)
@@ -2401,14 +2649,17 @@ const MapProject = () => {
}
aiFailuresInARowRef.current = 0
for (let index = 0; index < _rows.length; index++) {
- if (abortRef.current) break;
+ if (isStopped()) break;
// AI quota exhausted, or the AI service failing row after row: stop calling
// the AI step only. Matching itself already finished before this loop
// starts, so nothing else aborts.
if (shouldStopAIStep({quotaExhausted: aiQuotaExhaustedRef.current, failuresInARow: aiFailuresInARowRef.current})) break;
- await fetchRecommendation(_rows[index], resolvedPromptTemplate, true);
+ // A stopped run moves on; the row's call finishes in the background.
+ await untilCancelled(fetchRecommendation(_rows[index], resolvedPromptTemplate, true), isStopped);
}
+ if(!isCurrent())
+ return
const now = moment()
setBulkAIAnalysisEndedAt(now)
setEndMatchingAt(now)
@@ -2416,18 +2667,31 @@ const MapProject = () => {
setIsLoadingInDecisionView(false)
}
- const fetchBulkBridgeCandidates = async (_rows, algo) => {
+ // isStopped: the run's Stop, which stays true once a newer run starts.
+ const fetchBulkBridgeCandidates = async (_rows, algo, isStopped = () => abortRef.current) => {
+ // A row answered after a newer run started belongs to no run.
+ const runTicket = autoMatchRunTicketRef.current
setLoadingMatches(true)
setBridgeCandidatesStartedAt(moment())
for (let index = 0; index < _rows.length; index++) {
- if (abortRef.current) {
- setLoadingMatches(false)
+ if (isStopped()) {
+ // A newer run's loading state is its own.
+ if(runTicket === autoMatchRunTicketRef.current)
+ setLoadingMatches(false)
break;
};
if (matchQuotaStopRef.current) break;
+ if(isRunThrottled(algo.id)) {
+ markRowsThrottled(map(_rows.slice(index), '__index'), algo.id, getAlgoLogExtras(algo))
+ break
+ }
markAlgo(_rows[index].__index, algo.id, 0)
- await fetchBridgeCandidates(_rows[index], 0, undefined, undefined, undefined, false, true, ((response, payload) => {
+ // Stop doesn't wait on a request that never answers (ocl_issues#2849).
+ await untilCancelled(fetchBridgeCandidates(_rows[index], 0, undefined, undefined, undefined, false, true, ((response, payload) => {
+ if(runTicket !== autoMatchRunTicketRef.current)
+ return
+ saveLateRunResult()
const index = payload.rows[0].__index
const results = (isArray(response) ? response : response?.data)
log({action: 'algo_finished', extras: getAlgoLogExtras(algo)}, index)
@@ -2445,26 +2709,39 @@ const MapProject = () => {
rawResponse: response
}), {isRunTraffic: true})
}
- })); // wait for completion
+ })), isStopped); // wait for completion
await new Promise(resolve => setTimeout(resolve, 200)); // 1s delay
}
const now = moment()
setBridgeCandidatesEndedAt(now)
}
- const fetchBulkScispacyCandidates = async (_rows, algo) => {
+ // isStopped: the run's Stop, which stays true once a newer run starts.
+ const fetchBulkScispacyCandidates = async (_rows, algo, isStopped = () => abortRef.current) => {
+ // A row answered after a newer run started belongs to no run.
+ const runTicket = autoMatchRunTicketRef.current
setLoadingMatches(true)
setScispacyCandidatesStartedAt(moment())
for (let index = 0; index < _rows.length; index++) {
- if (abortRef.current) {
- setLoadingMatches(false)
+ if (isStopped()) {
+ // A newer run's loading state is its own.
+ if(runTicket === autoMatchRunTicketRef.current)
+ setLoadingMatches(false)
break;
};
+ if(isRunThrottled(algo.id) || isRunThrottled('ocl-scispacy-loinc')) {
+ markRowsThrottled(map(_rows.slice(index), '__index'), algo.id, getAlgoLogExtras(algo))
+ break
+ }
markAlgo(_rows[index].__index, algo.id, 0)
setLoadingMatches(true)
- await fetchScispacyCandidates(_rows[index], false, false, true, (response => {
+ // Stop doesn't wait on a request that never answers (ocl_issues#2849).
+ await untilCancelled(fetchScispacyCandidates(_rows[index], false, false, true, (response => {
+ if(runTicket !== autoMatchRunTicketRef.current)
+ return
+ saveLateRunResult()
const _index = _rows[index].__index
const results = [{row: _rows[index], results: fromScispacyResultsToConcepts(getScispacyRowResults(response.data, _index))}]
log({action: 'algo_finished', extras: getAlgoLogExtras(algo)}, _index)
@@ -2482,7 +2759,7 @@ const MapProject = () => {
rawResponse: response
}), {isRunTraffic: true})
}
- })); // wait for completion
+ })), isStopped); // wait for completion
await new Promise(resolve => setTimeout(resolve, 500)); // 1s delay
}
const now = moment()
@@ -3531,12 +3808,14 @@ const MapProject = () => {
// resolved, so a failed call showed as "no candidates" (ocl_online#283).
// requestSingleMatch retries it like a bulk batch; a call that still fails
// reaches the callback as an error body, which shows the error and marks
- // the algorithm failed for the row.
+ // the algorithm failed for the row. A 429 is waited out, and a server
+ // still busy at the cap marks the algorithm throttled (ocl_issues#2849).
requestSingleMatch(() => service.request(
'POST',
payload,
(algoDef.type === 'custom' && algoDef.url) ? algoDef.token : null,
{
+ ...(matchesOnOCL(algoDef) ? {headers: attrHeaders({rowIndex: __row.__index, algorithmId: algoDef.id, capacityAware: true})} : {}),
query: {
includeSearchMeta: true,
includeRetired: isBoolean(_retired) ? _retired : retired,
@@ -3549,11 +3828,19 @@ const MapProject = () => {
semantic: ['ocl-semantic', 'custom'].includes(algoDef.type),
reranker: !isMultiAlgo && algoDef.provider === 'ocl',
encoder_model: !isMultiAlgo && encoderModel ? encoderModel : undefined
- }
+ },
+ handlesThrottle: true,
+ timeout: HEAVY_REQUEST_TIMEOUT_MS,
}
- )).then(({ok, response, errorBody, previewLimit}) => {
+ ), {retryOptions: {
+ gate: getCapacityGate(service.URL),
+ onCapacity: noteCapacity,
+ ...trackCapacityWait([__row.__index], getAlgoLogExtras(algoDef)),
+ }}).then(({ok, response, errorBody, previewLimit, throttled}) => {
if(ok)
return callback(response, payload)
+ if(throttled)
+ return callback({detail: t('map_project.algorithm_throttled'), status: 429, throttled: true}, payload)
// A preview limit keeps the server's body, which opens the limit dialog.
callback(previewLimit ? errorBody : {...errorBody, detail: t('map_project.match_request_failed', {error: errorBody.detail})}, payload)
})
@@ -3618,6 +3905,18 @@ const MapProject = () => {
}
markAlgo(__row.__index, algoId, 0)
setIsLoadingInDecisionView(true)
+ // One algorithm failing doesn't stop the row's others, and the panel
+ // stops loading (ocl_online#283). ScispaCy and the bridge show their own
+ // failures, then go on here too (ocl_issues#2849).
+ const goOnAfterFailure = () => {
+ setIsLoadingInDecisionView(false)
+ const nextAlgoAfterFailure = getNextAlgoDef(algoId)
+ if(nextAlgoAfterFailure?.id && (offset === 0 || nextAlgoAfterFailure.type !== 'ocl-scispacy')) {
+ markAlgo(__row.__index, nextAlgoAfterFailure.id, 0)
+ fetchAllCandidatesForRow(nextAlgoAfterFailure.id, __row, offset, _retired, scrollToBottom, _filters, forceReload, true)
+ } else
+ scheduleRerank(__row.__index)
+ }
const onResponse = (response, payload) => {
const projectContext = buildProjectContext()
const logExtras = getAlgoLogExtras(getAlgoDef(algoId))
@@ -3626,13 +3925,18 @@ const MapProject = () => {
return
clearRefreshRowStageSnapshot(__row.__index)
refreshMapperQuotaCache()
- log({action: 'algo_failed', extras: {...logExtras, error: response.detail, status: response.status, ...(offset ? {offset} : {})}}, __row.__index)
- setAlert({message: response.detail, severity: 'error'})
+ // The server stayed too busy: the algorithm didn't run, it didn't
+ // fail, and no failed response is kept (ocl_issues#2849).
+ const isThrottled = Boolean(response.throttled)
+ log({action: isThrottled ? 'algo_throttled' : 'algo_failed', ...(isThrottled ? {description: t('map_project.row_throttled')} : {}), extras: {...logExtras, error: response.detail, status: response.status, ...(offset ? {offset} : {})}}, __row.__index)
+ setAlert({message: response.detail, severity: isThrottled ? 'warning' : 'error'})
if(offset) {
// A failed "load more" keeps the pages already loaded. The
// algorithm stays done if it loaded a page; "load more" also runs
// algorithms that never did, and those stay failed.
- markAlgo(__row.__index, algoId, hasSuccessfulAlgorithmResponse(rowMatchStateRef.current?.[__row.__index], algoId) ? 1 : -2)
+ markAlgo(__row.__index, algoId, hasSuccessfulAlgorithmResponse(rowMatchStateRef.current?.[__row.__index], algoId) ? 1 : (isThrottled ? -4 : -2))
+ } else if(isThrottled) {
+ markAlgo(__row.__index, algoId, -4)
} else {
markAlgo(__row.__index, algoId, -2)
mergeIntoRowMatchState(__row.__index, normalizeAlgorithmInvocation(null, {
@@ -3645,15 +3949,7 @@ const MapProject = () => {
rawResponse: response
}))
}
- // One algorithm failing doesn't stop the row's others, and the
- // panel stops loading (ocl_online#283).
- setIsLoadingInDecisionView(false)
- const nextAlgoAfterFailure = getNextAlgoDef(algoId)
- if(nextAlgoAfterFailure?.id && (offset === 0 || nextAlgoAfterFailure.type !== 'ocl-scispacy')) {
- markAlgo(__row.__index, nextAlgoAfterFailure.id, 0)
- fetchAllCandidatesForRow(nextAlgoAfterFailure.id, __row, offset, _retired, scrollToBottom, _filters, forceReload, true)
- } else
- scheduleRerank(__row.__index)
+ goOnAfterFailure()
return
}
log({action: 'algo_finished', extras: logExtras}, __row.__index)
@@ -3727,9 +4023,9 @@ const MapProject = () => {
if(['ocl-semantic', 'ocl-search', 'custom'].includes(algoDef.type)) {
fetchOCLOrCustomCandidates(algoDef, _row, offset, _retired, _filters, onResponse)
} else if (algoDef.type === 'ocl-scispacy') {
- fetchScispacyCandidates(__row, scrollToBottom, forceReload, false, onResponse)
+ fetchScispacyCandidates(__row, scrollToBottom, forceReload, false, onResponse, goOnAfterFailure)
} else if (['ocl-bridge', 'ocl-ciel-bridge'].includes(algoDef.type)) {
- fetchBridgeCandidates(__row, offset, _retired, scrollToBottom, _filters, forceReload, false, onResponse)
+ fetchBridgeCandidates(__row, offset, _retired, scrollToBottom, _filters, forceReload, false, onResponse, goOnAfterFailure)
}
} else {
setAlert({message: t('map_project.no_valid_columns_for_matching')})
@@ -3737,7 +4033,9 @@ const MapProject = () => {
}
}
- const fetchScispacyCandidates = async (_row, scrollToBottom, forceReload=false, isBulk=false, callback) => {
+ // onFailure: a row matched by hand goes on to its next algorithm after
+ // ScispaCy fails or stays throttled (ocl_issues#2849).
+ const fetchScispacyCandidates = async (_row, scrollToBottom, forceReload=false, isBulk=false, callback, onFailure) => {
let __row = isEmpty(_row) ? row : _row
// Gate on AlgorithmResponse presence, not candidate count: a successful
// invocation that returned zero matches still writes an algorithm_response
@@ -3780,73 +4078,112 @@ const MapProject = () => {
// the next row visit retries.
const SCISPACY_WARMUP_RETRY_MS = 2 * 60 * 1000
const SCISPACY_WARMUP_MAX_MS = 10 * 60 * 1000
+ const SCISPACY_ALGO_ID = 'ocl-scispacy-loinc'
+ const SCISPACY_ERROR_MESSAGE = "OCL's scispacy matching service is starting up. This may take a couple minutes. You can safely leave this row and come back. Click Refresh if results aren't here in a couple of minutes."
const warmupStart = Date.now()
- let warmingUp = true
let seenWarmingUp = false
+ // A run's waits end when it stops; a row matched by hand isn't stopped.
+ // (abortRef stays set after a Stop, so it can't gate a manual call.)
+ const isCancelled = isBulk ? getRunStopCheck() : () => false
+ // A run's call that ends after a newer run started leaves its stages alone.
+ const runTicket = autoMatchRunTicketRef.current
+ const isSuperseded = () => isBulk && runTicket !== autoMatchRunTicketRef.current
+ const capacityWait = trackCapacityWait([__row.__index], {algo: SCISPACY_ALGO_ID})
+ // A newer run's stages and panel state are its own.
+ const endStopped = () => {
+ if(isSuperseded())
+ return
+ markAlgo(__row.__index, SCISPACY_ALGO_ID, -1)
+ setIsLoadingInDecisionView(false)
+ }
- while(warmingUp) {
- if(abortRef.current) { setIsLoadingInDecisionView(false); return }
- try {
+ for(;;) {
+ // service.request rejects on a non-2xx: service.post resolved a 5xx as
+ // its body, and never settled on a 429 (ocl_issues#2849). A 429 waits
+ // for Retry-After; a network error or a gateway error is retried, but
+ // not a warm-up answer, which waits two minutes below.
+ const result = await requestWithCapacityRetry(() => service.request('POST', payload, null, {
// ocl_online#105 Phase 5: the scispacy service records its own
// api_transactions (ocl-scispacy-loinc) via its AnalyticsMiddleware, so
// attribute the $match like the others. Per-row (single-row payload) →
// scalar row_index; isBulk distinguishes a run call from a manual one.
- const response = await service.post(payload, null, attrHeaders({rowIndex: __row.__index, algorithmId: 'ocl-scispacy-loinc', isRunTraffic: isBulk}))
-
- // APIService resolves 5xx errors with error.response.data, so response
- // is the parsed body object — not the axios response. A 503 warming_up
- // from the lambda resolves as {status: 'warming_up', message: '...'}.
- // A 502 (EC2 still booting) resolves as {error: '...'}.
- const isWarmingUp = response?.status === 'warming_up' || (seenWarmingUp && response?.error)
- if(isWarmingUp) {
- seenWarmingUp = true
- const elapsed = Date.now() - warmupStart
- if(elapsed + SCISPACY_WARMUP_RETRY_MS > SCISPACY_WARMUP_MAX_MS) {
- markAlgo(__row.__index, 'ocl-scispacy-loinc', -2)
- log({action: 'algo_failed', extras: {algo: 'ocl-scispacy-loinc', status: 503, detail: 'warming_up_timeout'}}, __row.__index)
- setAlert({message: "OCL's scispacy service did not come up within 10 minutes.", severity: 'error'})
- setIsLoadingInDecisionView(false)
- return response
- }
- // Don't cover an error an earlier algorithm for this row showed.
- setAlert(prev => (prev?.severity === 'error' ? prev : {message: t('map_project.scispacy_warming_up'), severity: 'info'}))
- await new Promise(resolve => setTimeout(resolve, SCISPACY_WARMUP_RETRY_MS))
- } else {
- warmingUp = false
- const isError = response?.detail
- || response?.status >= 400
- || (response && response.data === undefined && response.status !== 200)
- if(isError) {
- if(handlePreviewLimitError(response, __row.__index, 'ocl-scispacy-loinc'))
- return response
- clearRefreshRowStageSnapshot(__row.__index)
- refreshMapperQuotaCache()
- markAlgo(__row.__index, 'ocl-scispacy-loinc', -2)
- log({action: 'algo_failed', extras: {algo: 'ocl-scispacy-loinc', status: response?.status, detail: response?.detail}}, __row.__index)
- setAlert({
- message: response?.detail || "OCL's scispacy matching service is starting up. This may take a couple minutes. You can safely leave this row and come back. Click Refresh if results aren't here in a couple of minutes.",
- severity: 'warning'
- })
- setIsLoadingInDecisionView(false)
- return response
- }
- // Clear ScispaCy's own warming-up or starting notice, but not an
- // error another algorithm for this row showed (ocl_online#283).
- setAlert(prev => (prev?.severity === 'error' ? prev : false))
- if(callback) callback(response, payload)
- setIsLoadingInDecisionView(false)
+ headers: attrHeaders({rowIndex: __row.__index, algorithmId: SCISPACY_ALGO_ID, isRunTraffic: isBulk}),
+ handlesThrottle: true,
+ timeout: HEAVY_REQUEST_TIMEOUT_MS,
+ }), {
+ gate: getCapacityGate(service.URL),
+ isCancelled,
+ maxWaitMs: isBulk ? CAPACITY_WAIT_CAP_MS : INTERACTIVE_WAIT_CAP_MS,
+ isRetryable: err => !isScispacyWarmingUp(err, seenWarmingUp) && isRetryableError(err),
+ onCapacity: noteCapacity,
+ ...capacityWait,
+ })
+
+ if(result.ok) {
+ const response = result.response
+ if(isSuperseded())
return response
- }
- } catch(err) {
- warmingUp = false
- markAlgo(__row.__index, 'ocl-scispacy-loinc', -2)
- log({action: 'algo_failed', extras: {algo: 'ocl-scispacy-loinc', error: err?.message}}, __row.__index)
- setAlert({
- message: "OCL's scispacy matching service is starting up. This may take a couple minutes. You can safely leave this row and come back. Click Refresh if results aren't here in a couple of minutes.",
- severity: 'warning'
- })
+ // Clear ScispaCy's own warming-up or starting notice, but not an
+ // error another algorithm for this row showed (ocl_online#283).
+ setAlert(prev => (prev?.severity === 'error' ? prev : false))
+ if(callback) callback(response, payload)
+ setIsLoadingInDecisionView(false)
+ return response
+ }
+ if(result.reason === 'cancelled') {
+ endStopped()
+ return
+ }
+ if(result.reason === 'throttled') {
+ if(isSuperseded())
+ return
+ // Not run, not failed: the row can be run again.
+ markAlgo(__row.__index, SCISPACY_ALGO_ID, -4)
+ log({action: 'algo_throttled', description: t('map_project.row_throttled'), extras: {algo: SCISPACY_ALGO_ID, attempts: result.attempts, waited_ms: result.waitedMs}}, __row.__index)
+ if(isBulk)
+ noteThrottledRunRow(SCISPACY_ALGO_ID, __row.__index)
+ setAlert(prev => (prev?.severity === 'error' ? prev : {message: t('map_project.algorithm_throttled'), severity: 'warning'}))
setIsLoadingInDecisionView(false)
+ onFailure?.()
+ return
}
+
+ const err = result.error
+ if(isSuperseded())
+ return
+ const body = err?.response?.data
+ // The service's host is booting: it answers 503 warming_up, then 502
+ // until the app is up.
+ if(isScispacyWarmingUp(err, seenWarmingUp)) {
+ seenWarmingUp = true
+ const elapsed = Date.now() - warmupStart
+ if(elapsed + SCISPACY_WARMUP_RETRY_MS > SCISPACY_WARMUP_MAX_MS) {
+ markAlgo(__row.__index, SCISPACY_ALGO_ID, -2)
+ log({action: 'algo_failed', extras: {algo: SCISPACY_ALGO_ID, status: 503, detail: 'warming_up_timeout'}}, __row.__index)
+ setAlert({message: "OCL's scispacy service did not come up within 10 minutes.", severity: 'error'})
+ setIsLoadingInDecisionView(false)
+ onFailure?.()
+ return body
+ }
+ // Don't cover an error an earlier algorithm for this row showed.
+ setAlert(prev => (prev?.severity === 'error' ? prev : {message: t('map_project.scispacy_warming_up'), severity: 'info'}))
+ if(!(await sleepUnlessCancelled(SCISPACY_WARMUP_RETRY_MS, isCancelled))) {
+ endStopped()
+ return
+ }
+ continue
+ }
+
+ if(handlePreviewLimitError(err, __row.__index, SCISPACY_ALGO_ID))
+ return body
+ clearRefreshRowStageSnapshot(__row.__index)
+ refreshMapperQuotaCache()
+ markAlgo(__row.__index, SCISPACY_ALGO_ID, -2)
+ log({action: 'algo_failed', extras: {algo: SCISPACY_ALGO_ID, status: err?.response?.status || null, detail: body?.detail, error: err?.message}}, __row.__index)
+ setAlert({message: body?.detail || SCISPACY_ERROR_MESSAGE, severity: 'warning'})
+ setIsLoadingInDecisionView(false)
+ onFailure?.()
+ return body
}
}
@@ -3924,10 +4261,33 @@ const MapProject = () => {
const rerank = async (_index, isBulk=false, isRunTraffic=isBulk) => {
const index = isNumber(_index) ? _index : rowIndex
if(!isNumber(index)) return null
- if(inFlightRerankRef.current.has(index)) {
- // Another rerank is in flight for this row; flag a rerun and bail.
- rerankRerunNeededRef.current.add(index)
- return null
+ // A run's rerank ends when the run stops, and a stopped run proposes no
+ // more mappings; a row reranked by hand isn't stopped (ocl_issues#2849).
+ const isStopped = isRunTraffic ? getRunStopCheck() : () => false
+ // The run this call belongs to, captured now: a newer run's stages are
+ // its own (ocl_issues#2849).
+ const runTicket = autoMatchRunTicketRef.current
+ const isSuperseded = () => isRunTraffic && runTicket !== autoMatchRunTicketRef.current
+ // A proposal lands a second later: skip it if the run stopped meanwhile.
+ const proposeMapping = () => setTimeout(() => {
+ if(!isStopped())
+ setAutoMatched([index])
+ }, 1000)
+ const rerankInFlight = inFlightRerankRef.current.get(index)
+ if(rerankInFlight) {
+ if(!isBulk) {
+ // Another rerank is in flight, or waiting its turn, for this row:
+ // flag a rerun for the candidates that arrive meanwhile.
+ rerankRerunNeededRef.current.add(index)
+ return null
+ }
+ // The run's sweep waits for it, then runs again for any candidates
+ // that arrived meanwhile, and proposes the row's mapping, which only
+ // the sweep's own call does (ocl_issues#2849).
+ await rerankInFlight
+ if(isStopped())
+ return null
+ return rerank(index, true, isRunTraffic)
}
// Wait for any in-flight $lookups for this row's concepts. buildRerankRowsForRow
// filters out nameless entries, so running before lookups complete silently
@@ -3940,6 +4300,13 @@ const MapProject = () => {
// Bounded, so a lookup that never settles can't stall rerank, and
// with it a run (ocl_online#283).
await waitForLookups(pendingLookups, RERANK_LOOKUP_WAIT_MS, stuckLookupsRef.current)
+ // The run may have stopped, and another started, during the wait:
+ // touch nothing of the new run's (ocl_issues#2849).
+ if(isStopped())
+ return null
+ // Another rerank for the row may have started meanwhile: join it.
+ if(inFlightRerankRef.current.has(index))
+ return rerank(index, isBulk, isRunTraffic)
// The settling lookups fired scheduleRerank via settle()/writeConceptCachePatch,
// setting a 300ms debounce timer. Clear it — we're already inside rerank() and
// about to run, so that timer would only queue a redundant rerankRerunNeeded.
@@ -3959,21 +4326,88 @@ const MapProject = () => {
// scheduleRerank that already scored the row from algo onResponse.
// Without the setAutoMatched trigger, auto-match would never propose
// a mapping even for rows with a clearly-recommended top candidate.
+ // Every candidate already has a score (e.g. inline, from a
+ // single-algorithm $match): the row's rerank is done.
+ const conceptRows = Object.values(rowMatchStateRef.current[index]?.concept_rows || {})
+ if(query && conceptRows.length && conceptRows.every(cr => isNumber(cr?.rerank_score)) && rowStageRef.current[index]?.rerank !== 1 && !isSuperseded())
+ markAlgo(index, 'rerank', 1)
if(isBulk && isNumber(index))
- setTimeout(() => setAutoMatched([index]), 1000)
+ proposeMapping()
+ return null
+ }
+ // A rerank in this run already stayed refused for the whole cap: don't
+ // ask again, the row's rerank is throttled ("not run, retry").
+ if(isRunTraffic && isRunThrottled('rerank')) {
+ markRowsThrottled([index], 'rerank', {algo: 'reranker'})
return null
}
- inFlightRerankRef.current.add(index)
+ let settleInFlight
+ const ownInFlight = new Promise(resolve => { settleInFlight = resolve })
+ inFlightRerankRef.current.set(index, ownInFlight)
markAlgo(index, 'rerank', 0)
const service = APIService.concepts().appendToUrl('$rerank/')
+ const isCancelled = isStopped
+ let release = null
try {
- const response = await service.post({
+ // At most RERANK_MAX_IN_FLIGHT $rerank calls at once, whatever fired
+ // them: every finished batch used to fire one per row (ocl_issues#2849).
+ release = await rerankLimiterRef.current.acquire({isCancelled})
+ if(!release) {
+ // Stopped before its turn; a stopped run reranks nothing more.
+ rerankRerunNeededRef.current.delete(index)
+ if(!isSuperseded())
+ markAlgo(index, 'rerank', -1)
+ return null
+ }
+ // Another of the run's reranks stayed refused for the whole cap while
+ // this one waited its turn: don't ask again.
+ if(isRunTraffic && isRunThrottled('rerank')) {
+ markRowsThrottled([index], 'rerank', {algo: 'reranker'})
+ return null
+ }
+ // Score the candidates that arrived while this call waited its turn too.
+ const latestRerankRows = buildRerankRowsForRow(index)
+ // service.request rejects on a non-2xx; service.post resolved a 5xx as if
+ // it scored nothing, and never settled on a 429. A 429 is waited out.
+ const result = await requestWithCapacityRetry(() => service.request('POST', {
q: query,
- rows: rerankRows,
+ rows: latestRerankRows.length ? latestRerankRows : rerankRows,
...(encoderModel ? { encoder_model: encoderModel } : {})
- }, null, attrHeaders({rowIndex: index, algorithmId: 'reranker', isRunTraffic}))
- if(handlePreviewLimitError(response, index, 'rerank', {isRun: isRunTraffic, stopPhase: 'rerank'}))
- return response
+ }, null, {headers: attrHeaders({rowIndex: index, algorithmId: 'reranker', isRunTraffic, capacityAware: true}), handlesThrottle: true, timeout: HEAVY_REQUEST_TIMEOUT_MS}), {
+ gate: getCapacityGate(service.URL),
+ isCancelled,
+ onCapacity: noteCapacity,
+ ...trackCapacityWait([index], {algo: 'reranker'}),
+ })
+ if(result.reason === 'cancelled') {
+ rerankRerunNeededRef.current.delete(index)
+ if(!isSuperseded())
+ markAlgo(index, 'rerank', -1)
+ return null
+ }
+ if(result.reason === 'throttled') {
+ log({action: 'rerank_throttled', description: t('map_project.row_throttled'), extras: {attempts: result.attempts, waited_ms: result.waitedMs}}, index)
+ if(isSuperseded())
+ return null
+ markAlgo(index, 'rerank', -4)
+ if(isRunTraffic)
+ noteThrottledRunRow('rerank', index)
+ return null
+ }
+ if(!result.ok) {
+ if(isSuperseded())
+ return null
+ if(handlePreviewLimitError(result.error, index, 'rerank', {isRun: isRunTraffic, stopPhase: 'rerank'}))
+ return null
+ log({action: 'rerank_failed', description: `Rerank failed with ${encoderModel}`, extras: {status: result.error?.response?.status || null, error: result.error?.response?.data?.detail || result.error?.message || null}}, index)
+ markAlgo(index, 'rerank', -2)
+ return null
+ }
+ const response = result.response
+ // A newer run may rank with another encoder or input: its scores are
+ // its own (ocl_issues#2849).
+ if(isSuperseded())
+ return null
// Write rerank_score into the row's ConceptRows. matchRerankResultToKey
// throws on canonical-identity miss; surface to the alert state so a
@@ -4011,23 +4445,30 @@ const MapProject = () => {
// (Removed legacy allCandidates re-stamp: rerank_score now lives on
// rowMatchState[idx].concept_rows[concept_key], which we updated above.)
- markAlgo(index, 'rerank', 1)
+ // The scores stand either way; a newer run's stages are its own.
+ if(!isSuperseded())
+ markAlgo(index, 'rerank', 1)
+ if(isRunTraffic)
+ saveLateRunResult()
log({action: 'rerank_finished', description: `Reranked with ${encoderModel}`}, index)
if(isBulk)
- setTimeout(() => setAutoMatched([index]), 1000)
+ proposeMapping()
return response
} catch (e) {
- if(handlePreviewLimitError(e, index, 'rerank', {isRun: isRunTraffic, stopPhase: 'rerank'}))
- return null
- log({action: 'rerank_failed', description: `Rerank failed with ${encoderModel}`}, index)
+ log({action: 'rerank_failed', description: `Rerank failed with ${encoderModel}`, extras: {error: e?.message || null}}, index)
markAlgo(index, 'rerank', -2)
return null
} finally {
- inFlightRerankRef.current.delete(index)
- // If new ConceptRows arrived while we were in flight, fire again.
+ release?.()
+ if(inFlightRerankRef.current.get(index) === ownInFlight)
+ inFlightRerankRef.current.delete(index)
+ settleInFlight()
+ // If new ConceptRows arrived while we were in flight, fire again,
+ // through the current render's scheduler: this call's closure may hold
+ // an older encoder or row input (ocl_issues#2849).
if(rerankRerunNeededRef.current.has(index)) {
rerankRerunNeededRef.current.delete(index)
- scheduleRerank(index)
+ scheduleRerankRef.current(index)
}
}
}
@@ -4060,15 +4501,19 @@ const MapProject = () => {
// to mapper-ui-manual. Internally consistent: if we treat this as a
// bulk-run rerank for gating, we attribute it to the run.
const isRunTraffic = Boolean(isBulkMatchRunningRef.current && bulkMatchAlgoIdsRef.current.length > 0)
+ // An algorithm the server stayed too busy for (-4), or that isn't
+ // available to this user (-3), is settled too: the row's other
+ // candidates still get ranked (ocl_issues#2849).
+ const isSettled = value => value === 1 || value === -2 || value === -3 || value === -4
if(isRunTraffic) {
// Bulk auto-match: wait until all selected algos are done (success or failed) for this row.
- const allDone = bulkMatchAlgoIdsRef.current.every(id => stage[id] === 1 || stage[id] === -2)
+ const allDone = bulkMatchAlgoIdsRef.current.every(id => isSettled(stage[id]))
if(!allDone) return
} else {
// Single-row: wait until every configured algo has either completed (1) or failed (-2).
// Checking only === 0 (in-flight) wasn't enough — algos not yet started (-1) or
// uninitialized (undefined) would let rerank fire before their candidates arrived.
- if(selectedAlgoIds.some(id => stage[id] !== 1 && stage[id] !== -2)) return
+ if(selectedAlgoIds.some(id => !isSettled(stage[id]))) return
}
const rowState = rowMatchStateRef.current[rowIndex]
if(!rowState) return
@@ -4087,8 +4532,14 @@ const MapProject = () => {
if(!hasEligiblePending) return
if(rerankDebounceRef.current[rowIndex])
clearTimeout(rerankDebounceRef.current[rowIndex])
+ // A run's rerank belongs to that run: if a newer run started before the
+ // timer fired, drop it (the newer run reranks its own rows) rather than
+ // let it pass as the newer run's (ocl_issues#2849).
+ const scheduledTicket = autoMatchRunTicketRef.current
rerankDebounceRef.current[rowIndex] = setTimeout(() => {
delete rerankDebounceRef.current[rowIndex]
+ if(isRunTraffic && scheduledTicket !== autoMatchRunTicketRef.current)
+ return
// isBulk=false (don't double the setAutoMatched side-effect the explicit
// bulk rerank already performs); isRunTraffic carries the run attribution.
rerank(rowIndex, false, isRunTraffic)
@@ -4222,7 +4673,9 @@ const MapProject = () => {
return formatted
}
- const fetchBridgeCandidates = (_row, offset=0, _retired, scrollToBottom, _filters, forceReload=false, isBulk=false, callback) => {
+ // onFailure: a row matched by hand goes on to its next algorithm after the
+ // bridge fails or stays throttled (ocl_issues#2849).
+ const fetchBridgeCandidates = (_row, offset=0, _retired, scrollToBottom, _filters, forceReload=false, isBulk=false, callback, onFailure) => {
let __row = isEmpty(_row) ? row : _row
const bridgeAlgoId = bridgeAlgo?.id || 'ocl-ciel-bridge'
const rowEntry = rowMatchStateRef.current?.[__row.__index]
@@ -4241,46 +4694,104 @@ const MapProject = () => {
setIsLoadingInDecisionView(true)
const payload = getPayloadForMatching([__row], repo)
let __offset = offset || 0
+ const failBridgeRow = message => {
+ clearRefreshRowStageSnapshot(__row.__index)
+ markAlgo(__row.__index, bridgeAlgoId, -2)
+ log({action: 'algo_failed', extras: {...getAlgoLogExtras(bridgeAlgo), error: message}}, __row.__index)
+ setAlert({message, severity: 'error'})
+ setIsLoadingInDecisionView(false)
+ }
return new Promise(resolve => {
- bridgeRef.current?.fetchBridgeCandidates(
- payload,
- __offset,
- isBoolean(_retired) ? _retired : retired,
- (candidates) => {
- if(callback) {
- const newCandidates = candidates?.map(candidate => ({
- ...candidate,
- results: candidate?.results?.map(result => ({
- ...result,
- search_meta: {
- ...result?.search_meta,
- algorithm: bridgeAlgoId
- }
- }))
- }));
- callback(newCandidates, payload)
- }
- resolve()
- },
- (response, errorMsg) => {
- if(handlePreviewLimitError(response, __row.__index, bridgeAlgoId)) {
+ // The component calls its service's post() synchronously from here; the
+ // context tells that post() the row, and whether a Stop ends its waits
+ // (ocl_issues#2849).
+ const context = {rowIndex: __row.__index, isBulk, isCancelled: isBulk ? getRunStopCheck() : () => false, posted: false}
+ // A run's call that ends after a newer run started leaves its stages alone.
+ const runTicket = autoMatchRunTicketRef.current
+ const isSuperseded = () => isBulk && runTicket !== autoMatchRunTicketRef.current
+ bridgeCallContextRef.current = context
+ try {
+ bridgeRef.current?.fetchBridgeCandidates(
+ payload,
+ __offset,
+ isBoolean(_retired) ? _retired : retired,
+ (candidates) => {
+ if(callback) {
+ const newCandidates = candidates?.map(candidate => ({
+ ...candidate,
+ results: candidate?.results?.map(result => ({
+ ...result,
+ search_meta: {
+ ...result?.search_meta,
+ algorithm: bridgeAlgoId
+ }
+ }))
+ }));
+ callback(newCandidates, payload)
+ }
resolve()
- return
- }
- clearRefreshRowStageSnapshot(__row.__index)
- refreshMapperQuotaCache()
- markAlgo(__row.__index, bridgeAlgoId, -2)
- log({action: 'algo_failed', extras: getAlgoLogExtras(bridgeAlgo)}, __row.__index)
- setAlert({message: response?.detail || errorMsg, severity: 'error'})
- setIsLoadingInDecisionView(false)
- resolve()
- },
- // ocl_online#105 Phase 5: attribution headers for the bridge $match,
- // applied by the premium component (react-bridge-match). Per-row here
- // (single-row payload) → scalar row_index. isBulk distinguishes a run
- // call from a manual per-row bridge fetch.
- attrHeaders({rowIndex: __row.__index, algorithmId: bridgeAlgoId, isRunTraffic: isBulk})
- )
+ },
+ (response, errorMsg) => {
+ if(isSuperseded()) {
+ resolve()
+ return
+ }
+ if(response?.cancelled) {
+ markAlgo(__row.__index, bridgeAlgoId, -1)
+ setIsLoadingInDecisionView(false)
+ resolve()
+ return
+ }
+ // The server stayed too busy: not run, not failed.
+ if(response?.throttled) {
+ clearRefreshRowStageSnapshot(__row.__index)
+ markAlgo(__row.__index, bridgeAlgoId, -4)
+ log({action: 'algo_throttled', description: t('map_project.row_throttled'), extras: getAlgoLogExtras(bridgeAlgo)}, __row.__index)
+ if(isBulk)
+ noteThrottledRunRow(bridgeAlgoId, __row.__index)
+ setAlert(prev => (prev?.severity === 'error' ? prev : {message: response.detail, severity: 'warning'}))
+ setIsLoadingInDecisionView(false)
+ onFailure?.()
+ resolve()
+ return
+ }
+ if(handlePreviewLimitError(response, __row.__index, bridgeAlgoId)) {
+ resolve()
+ return
+ }
+ clearRefreshRowStageSnapshot(__row.__index)
+ refreshMapperQuotaCache()
+ markAlgo(__row.__index, bridgeAlgoId, -2)
+ log({action: 'algo_failed', extras: getAlgoLogExtras(bridgeAlgo)}, __row.__index)
+ setAlert({message: response?.detail || errorMsg, severity: 'error'})
+ setIsLoadingInDecisionView(false)
+ onFailure?.()
+ resolve()
+ },
+ // ocl_online#105 Phase 5: attribution headers for the bridge $match,
+ // applied by the premium component (react-bridge-match). Per-row here
+ // (single-row payload) → scalar row_index. isBulk distinguishes a run
+ // call from a manual per-row bridge fetch.
+ attrHeaders({rowIndex: __row.__index, algorithmId: bridgeAlgoId, isRunTraffic: isBulk, capacityAware: true})
+ )
+ } catch (err) {
+ failBridgeRow(err?.message || t('unknown_error'))
+ onFailure?.()
+ resolve()
+ return
+ } finally {
+ bridgeCallContextRef.current = null
+ }
+ // The component sent nothing (bridging isn't available to this user or
+ // project), so no callback will come: mark it n/a and settle, rather
+ // than leave the row "Running" and its rerank waiting for good.
+ if(!context.posted) {
+ markAlgo(__row.__index, bridgeAlgoId, -3)
+ log({action: 'algo_skipped', extras: {...getAlgoLogExtras(bridgeAlgo), reason: 'bridge_unavailable'}}, __row.__index)
+ setIsLoadingInDecisionView(false)
+ onFailure?.()
+ resolve()
+ }
})
}
@@ -4426,8 +4937,15 @@ const MapProject = () => {
} else {
service = service.overrideURL(oclUrl)
}
- const response = await service.request('GET', undefined, authToken, {query: {includeMappings: true, mappingBrief: true, mapTypes: 'SAME-AS,SAME AS,SAME_AS', verbose: true}, headers: attrHeaders({rowIndex: logRowIndex, isRunTraffic})})
- const data = response?.data
+ // A busy server's 429 is waited out, and a gateway error retried
+ // (ocl_issues#2849); a 404 is retried below.
+ const result = await requestWithCapacityRetry(
+ () => service.request('GET', undefined, authToken, {query: {includeMappings: true, mappingBrief: true, mapTypes: 'SAME-AS,SAME AS,SAME_AS', verbose: true}, headers: attrHeaders({rowIndex: logRowIndex, isRunTraffic}), handlesThrottle: true, timeout: LIGHT_REQUEST_TIMEOUT_MS}),
+ {maxWaitMs: INTERACTIVE_WAIT_CAP_MS, onCapacity: noteCapacity}
+ )
+ if(!result.ok)
+ throw result.error || new Error(result.reason)
+ const data = result.response?.data
if(data?.id) return data
urlFetchCacheRef.current.delete(oclUrl)
return null
@@ -4481,15 +4999,24 @@ const MapProject = () => {
const body = toResolve.map(({reference}) => reference.version
? {url: reference.url, version: reference.version}
: {url: reference.url})
- // service.request rejects on any error, a 429 included, so the catch
- // below marks these lookups failed and settles them. APIService.post
- // answers a 429 with a promise that never settles, which left every
- // wait on these lookups hanging (ocl_online#283).
- resolvePromise = APIService.new()
- .overrideURL('/$resolveReference/')
- .request('POST', body, currentToken, {
+ // A busy server's 429 is waited out (ocl_issues#2849); any other error,
+ // or a server still busy at the cap, lands in the catch below, which
+ // marks these lookups failed and settles them. APIService.post answers
+ // a 429 with a promise that never settles, which left every wait on
+ // these lookups hanging (ocl_online#283).
+ resolvePromise = requestWithCapacityRetry(
+ () => APIService.new().overrideURL('/$resolveReference/').request('POST', body, currentToken, {
headers: attrHeaders({rowIndex: logRowIndex, isRunTraffic}),
query: resolveNamespace ? {namespace: resolveNamespace} : {},
+ handlesThrottle: true,
+ timeout: LIGHT_REQUEST_TIMEOUT_MS,
+ }),
+ {maxWaitMs: INTERACTIVE_WAIT_CAP_MS, onCapacity: noteCapacity}
+ )
+ .then(result => {
+ if(!result.ok)
+ throw result.error || new Error(result.reason)
+ return result.response
})
.then(async response => {
const items = Array.isArray(response?.data) ? response.data : []
@@ -4574,7 +5101,7 @@ const MapProject = () => {
has_filters: !isEmpty(appliedFilters || appliedFacets[rowIndex])
})
setIsLoadingInDecisionView(true)
- getLookupService().get(lookupConfig?.token, null, {
+ const query = {
includeSearchMeta: true,
includeMappings: true,
mappingBrief: true,
@@ -4585,8 +5112,24 @@ const MapProject = () => {
page: page || 1,
includeRetired: includeRetired === undefined ? retired : includeRetired,
...getFacetQueryParam(appliedFilters || appliedFacets[rowIndex]),
- }).then(response => {
- let items = response.data
+ }
+ // service.get never settled on a 429 and resolved an error as its body,
+ // which then threw here. A 429 is waited out; an error says so
+ // (ocl_issues#2849).
+ requestWithCapacityRetry(
+ () => getLookupService().request('GET', undefined, lookupConfig?.token, {query, handlesThrottle: true, timeout: LIGHT_REQUEST_TIMEOUT_MS}),
+ {maxWaitMs: INTERACTIVE_WAIT_CAP_MS, onCapacity: noteCapacity, ...trackCapacityWait([])}
+ ).then(result => {
+ if(!result.ok) {
+ setIsLoadingInDecisionView(false)
+ if(result.reason === 'throttled')
+ setAlert({message: t('map_project.search_throttled'), severity: 'warning'})
+ else
+ setAlert({message: t('map_project.search_failed', {error: result.error?.response?.data?.detail || result.error?.message || t('unknown_error')}), severity: 'error'})
+ return
+ }
+ const response = result.response
+ let items = response.data || []
setSearchedConcepts({...searchedConcepts, [row.__index]: items})
setSearchResponse(response)
setIsLoadingInDecisionView(false)
@@ -4632,7 +5175,15 @@ const MapProject = () => {
})
}
- const request = getLookupService().get(lookupConfig?.token, null, query).then(response => {
+ // A failed facets call isn't cached, so the next look tries again
+ // (ocl_issues#2849).
+ const request = requestWithCapacityRetry(
+ () => getLookupService().request('GET', undefined, lookupConfig?.token, {query, handlesThrottle: true, timeout: LIGHT_REQUEST_TIMEOUT_MS}),
+ {maxWaitMs: INTERACTIVE_WAIT_CAP_MS, onCapacity: noteCapacity, ...trackCapacityWait([])}
+ ).then(result => {
+ if(!result.ok)
+ return undefined
+ const response = result.response
const fields = response?.data?.facets?.fields || {}
setFetchedFacets(prev => ({...prev, [requestKey]: fields}))
if(latestFacetRequestRef.current[targetRowIndex] === requestKey)
@@ -5309,7 +5860,7 @@ const MapProject = () => {
bridgeEnabled={bridgeEnabled}
canBridge={canBridge}
isCoreUser={isCoreUser}
- fullRequestLimits={fullRequestLimits}
+ requestLimits={requestLimits}
canSelectAIModel={canSelectAIModel}
canSetAIOutputLocale={canSetAIOutputLocale}
canScispacy={canScispacy}
@@ -5373,7 +5924,7 @@ const MapProject = () => {
const mappedSrcs = bridgeMappedSources[bridgeUrl] || []
return Boolean(repoVersion?.url) && Boolean(bridgeUrl) && mappedSrcs.length > 0 &&
{
{abortRef.current ? t('map_project.stopping_gracefully') : t('map_project.stop_processing')}
}
+ {
+ // 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 &&
+
+ }
+ color='warning'
+ variant='outlined'
+ size='small'
+ label={isSplitView ? t('map_project.waiting_for_capacity_short') : t('map_project.waiting_for_capacity')}
+ sx={{margin: '5px'}}
+ />
+
+ }
}
{
@@ -5899,6 +6466,7 @@ const MapProject = () => {
candidatesScore={candidatesScore}
rowIndex={rowIndex}
rowStage={rowStageRef.current[rowIndex]}
+ capacityWait={Boolean(capacityWaitRows?.[rowIndex])}
rowState={rowMatchStateRef.current[rowIndex]}
conceptCache={conceptCache}
targetCanonical={buildProjectContext()?.target_repo?.canonical_url}
diff --git a/src/components/map-projects/MultiAlgoSelector.jsx b/src/components/map-projects/MultiAlgoSelector.jsx
index 105e9b3..9ce83c3 100644
--- a/src/components/map-projects/MultiAlgoSelector.jsx
+++ b/src/components/map-projects/MultiAlgoSelector.jsx
@@ -111,12 +111,11 @@ export default function MultiAlgoSelector({
maxAlgos=5,
repo,
isCoreUser,
- fullRequestLimits=false,
+ requestLimits=getRequestLimits(null),
}) {
const { t } = useTranslation()
- // Non-core users get at most 10 rows per batch and 5 concurrent requests
- // (ocl_online#274).
- const requestLimits = getRequestLimits(fullRequestLimits)
+ // Non-core users get at most 10 rows per batch, and early-access users 5
+ // concurrent requests, preview users 2 (ocl_online#274, ocl_issues#2849).
const { setAlert } = React.useContext(OperationsContext)
const [expanded, setExpanded] = useState(() => new Map());
const [errors, setErrors] = React.useState({})
@@ -186,8 +185,8 @@ export default function MultiAlgoSelector({
const getRequestFieldsSettings = (sel, algo) => getRequestSettings({
batch_size: sel.batch_size ?? algo?.batch_size,
concurrent_requests: sel.concurrent_requests ?? algo?.concurrent_requests,
- }, fullRequestLimits)
- const getLimitHelperText = max => fullRequestLimits ? undefined : t('map_project.request_limit_up_to', {max})
+ }, requestLimits)
+ const getLimitHelperText = max => requestLimits.capped ? t('map_project.request_limit_up_to', {max}) : undefined
const getBatchSizeProps = (sel, algo) => ({
value: getRequestFieldsSettings(sel, algo).batchSize,
onChange: e => updateSelected(sel.__key, { batch_size: clampInt(e.target.value, 1, requestLimits.batchSize) }),
@@ -314,8 +313,8 @@ export default function MultiAlgoSelector({
...omit(algo, ['getIcon', 'disabled', 'description', 'url']),
id: id,
name: name,
- batch_size: getRequestSettings(algo, fullRequestLimits).batchSize,
- concurrent_requests: getRequestSettings(algo, fullRequestLimits).concurrentRequests,
+ batch_size: getRequestSettings(algo, requestLimits).batchSize,
+ concurrent_requests: getRequestSettings(algo, requestLimits).concurrentRequests,
__key: Math.random(100).toString()
};
diff --git a/src/components/map-projects/__tests__/autoMatchRows.test.js b/src/components/map-projects/__tests__/autoMatchRows.test.js
index 32b1d78..64e0694 100644
--- a/src/components/map-projects/__tests__/autoMatchRows.test.js
+++ b/src/components/map-projects/__tests__/autoMatchRows.test.js
@@ -3,8 +3,9 @@ import assert from 'node:assert/strict'
import {
getPreviewEligibleRowIndexes, getRowsToProcess, spendsMatchQuota, getRowCapByMatchOperations, shouldStopAIStep,
- getAIRequestIdempotencyKey, getCandidatePoolFingerprint, hasCurrentAnalysis, getScispacyRowResults, getPendingRowLookups, waitForLookups, RERANK_LOOKUP_WAIT_MS, AI_LOOKUP_WAIT_MS
+ getAIRequestIdempotencyKey, getCandidatePoolFingerprint, hasCurrentAnalysis, getScispacyRowResults, getPendingRowLookups, waitForLookups, RERANK_LOOKUP_WAIT_MS, AI_LOOKUP_WAIT_MS, isScispacyWarmingUp, RERANK_MAX_IN_FLIGHT
} from '../autoMatchRows.js'
+import { createLimiter } from '../../../services/capacity.js'
const rows = [
{ __index: 0, label: 'zero' },
@@ -289,3 +290,43 @@ test('RERANK_LOOKUP_WAIT_MS: a hang guard, longer than the AI step\'s wait', ()
assert.ok(RERANK_LOOKUP_WAIT_MS >= 30000)
assert.ok(RERANK_LOOKUP_WAIT_MS > AI_LOOKUP_WAIT_MS)
})
+
+// ocl_issues#2849: ScispaCy now goes through service.request, which rejects on
+// a non-2xx. Its warm-up answers must still read as "wait", not as a failure.
+test('isScispacyWarmingUp: a 503 warming_up body means the service is starting', () => {
+ const err = {response: {status: 503, data: {status: 'warming_up', message: 'Starting'}}}
+ assert.equal(isScispacyWarmingUp(err), true)
+})
+
+test('isScispacyWarmingUp: a 502 while the host boots counts only once warm-up was seen', () => {
+ const err = {response: {status: 502, data: {error: 'Bad gateway'}}}
+ assert.equal(isScispacyWarmingUp(err), false)
+ assert.equal(isScispacyWarmingUp(err, true), true)
+})
+
+test('isScispacyWarmingUp: other errors are not warm-up', () => {
+ assert.equal(isScispacyWarmingUp({response: {status: 500, data: {detail: 'boom'}}}, true), false)
+ assert.equal(isScispacyWarmingUp({response: {status: 429, data: {}}}, true), false)
+ assert.equal(isScispacyWarmingUp(new Error('Network Error'), true), false)
+ assert.equal(isScispacyWarmingUp(undefined), false)
+})
+
+// ocl_issues#2849: every $rerank goes through one limiter per tab. A finished
+// 10-row batch used to fire 10 at once, on top of the rerank sweep's 2.
+test('RERANK_MAX_IN_FLIGHT: a batch\'s reranks and the sweep\'s together keep at most 2 in flight', async () => {
+ assert.equal(RERANK_MAX_IN_FLIGHT, 2)
+ const limiter = createLimiter(RERANK_MAX_IN_FLIGHT)
+ let inFlight = 0
+ let maxInFlight = 0
+ const rerank = async () => {
+ const release = await limiter.acquire()
+ inFlight += 1
+ maxInFlight = Math.max(maxInFlight, inFlight)
+ await new Promise(resolve => setTimeout(resolve, 2))
+ inFlight -= 1
+ release()
+ }
+ // ten per-row reranks as one batch finishes, and the sweep's two
+ await Promise.all([...Array(10).keys(), 'sweep-a', 'sweep-b'].map(rerank))
+ assert.equal(maxInFlight, 2)
+})
diff --git a/src/components/map-projects/__tests__/autosave.test.js b/src/components/map-projects/__tests__/autosave.test.js
index 78b9477..181c180 100644
--- a/src/components/map-projects/__tests__/autosave.test.js
+++ b/src/components/map-projects/__tests__/autosave.test.js
@@ -13,7 +13,7 @@
import test from 'node:test'
import assert from 'node:assert/strict'
-import { createAutosaveScheduler, saveOnLeave, trackSave, whenSaved, hasSaveInFlight, AUTOSAVE_DELAY_MS } from '../autosave.js'
+import { createAutosaveScheduler, createLatestSender, saveOnLeave, trackSave, whenSaved, hasSaveInFlight, AUTOSAVE_DELAY_MS } from '../autosave.js'
// Virtual clock standing in for setTimeout/clearTimeout. Timers scheduled
// while the clock is being advanced land in a later window, as real timers
@@ -304,3 +304,92 @@ test('trackSave: an older save settling does not untrack a newer one for the sam
await nextTick()
assert.equal(hasSaveInFlight(), false)
})
+
+// ── createLatestSender (ocl_issues#2849) ───────────────────────────────────
+// The logs POST sends the whole log each time. With retries, an older POST
+// that waited out a 429 could land after a newer one and overwrite it; so one
+// POST at a time, and triggers during it coalesce into one more, which reads
+// the log as it is then.
+
+const deferred = () => {
+ let resolve
+ const promise = new Promise(r => { resolve = r })
+ return {promise, resolve}
+}
+
+test('createLatestSender: one trigger sends once', async () => {
+ let sends = 0
+ const sender = createLatestSender(async () => { sends += 1 })
+ await sender.trigger()
+ assert.equal(sends, 1)
+})
+
+test('createLatestSender: triggers while a send is in flight coalesce into one more send', async () => {
+ const gates = [deferred(), deferred()]
+ let sends = 0
+ const sender = createLatestSender(() => gates[sends++].promise)
+ const first = sender.trigger()
+ sender.trigger()
+ sender.trigger()
+ const last = sender.trigger()
+ assert.equal(sends, 1)
+ gates[0].resolve()
+ await new Promise(resolve => setTimeout(resolve, 0))
+ assert.equal(sends, 2)
+ gates[1].resolve()
+ // every trigger's promise settles once the log is sent as it stood last
+ await Promise.all([first, last])
+ assert.equal(sends, 2)
+})
+
+test('createLatestSender: a failed send doesn\'t stop the pending one, and nothing rejects', async () => {
+ let sends = 0
+ const sender = createLatestSender(async () => {
+ sends += 1
+ if(sends === 1) {
+ await new Promise(resolve => setTimeout(resolve, 5))
+ throw new Error('boom')
+ }
+ })
+ const first = sender.trigger()
+ const second = sender.trigger()
+ await Promise.all([first, second])
+ assert.equal(sends, 2)
+})
+
+test('createLatestSender: after it settles, the next trigger sends again', async () => {
+ let sends = 0
+ const sender = createLatestSender(async () => { sends += 1 })
+ await sender.trigger()
+ await sender.trigger()
+ assert.equal(sends, 2)
+})
+
+// Codex review, pass 1: the payload belongs to the sender (one per project),
+// so a retry for project A can't pick up project B's log from shared state.
+test('createLatestSender: the send reads the latest payload triggered on this sender, on every attempt', async () => {
+ const gate = deferred()
+ const seen = []
+ const sender = createLatestSender(async getPayload => {
+ seen.push(getPayload())
+ if(seen.length === 1) {
+ await gate.promise
+ // a retry within the same send reads it again
+ seen.push(getPayload())
+ }
+ })
+ const first = sender.trigger({logs: 'a1'})
+ sender.trigger({logs: 'a2'})
+ sender.trigger({logs: 'a3'})
+ gate.resolve()
+ await first
+ assert.deepEqual(seen, [{logs: 'a1'}, {logs: 'a3'}, {logs: 'a3'}])
+})
+
+test('createLatestSender: two senders keep their own payloads', async () => {
+ const seen = []
+ const a = createLatestSender(async getPayload => { seen.push(['a', getPayload()]) })
+ const b = createLatestSender(async getPayload => { seen.push(['b', getPayload()]) })
+ await Promise.all([a.trigger('project A'), b.trigger('project B')])
+ assert.deepEqual(seen.sort(), [['a', 'project A'], ['b', 'project B']])
+})
diff --git a/src/components/map-projects/__tests__/matchBatch.test.js b/src/components/map-projects/__tests__/matchBatch.test.js
index 505d3dd..4d4a24b 100644
--- a/src/components/map-projects/__tests__/matchBatch.test.js
+++ b/src/components/map-projects/__tests__/matchBatch.test.js
@@ -12,7 +12,8 @@
import test from 'node:test'
import assert from 'node:assert/strict'
-import { formatRowNumbers, isRetryableMatchError, runMatchBatch, requestSingleMatch } from '../matchBatch.js'
+import { formatRowNumbers, isRetryableMatchError, runMatchBatch, requestSingleMatch, runWithConcurrency } from '../matchBatch.js'
+import { createCapacityGate } from '../../../services/capacity.js'
import { createLatestRequestGate } from '../aiVisibility.js'
const NO_WAIT = {baseDelayMs: 0}
@@ -286,3 +287,268 @@ test('requestSingleMatch: a preview-limit 403 passes its body through unchanged
assert.equal(result.previewLimit, true)
assert.deepEqual(result.errorBody, limit)
})
+
+// ── Waiting out a busy server (OpenConceptLab/ocl_issues#2849) ────────────────
+// A 429 used to fail the batch at once. Now it waits for Retry-After (plus
+// jitter) and resumes; a batch still throttled after the long cap is marked
+// throttled (-4, "not run, retry"), never failed.
+
+const throttled = () => httpError(429, {})
+const withRetryAfter = (err, seconds) => { err.response.headers = {'retry-after': String(seconds)}; return err }
+const virtualClock = () => {
+ const clock = {t: 0}
+ clock.now = () => clock.t
+ clock.sleep = async ms => { clock.t += ms }
+ return clock
+}
+
+test('runMatchBatch: a 429 waits for Retry-After and the batch then succeeds; its rows are never marked failed', async () => {
+ const rec = recorder()
+ const clock = virtualClock()
+ const waits = []
+ let calls = 0
+ const send = async () => {
+ calls += 1
+ if(calls === 1) throw withRetryAfter(throttled(), 12)
+ return ok([{row: {__index: 0}, results: []}, {row: {__index: 1}, results: []}])
+ }
+
+ const data = await runMatchBatch({
+ rowIndexes: [0, 1], send, ...rec,
+ onWait: info => waits.push(info.reason), onWaitEnd: () => waits.push('end'),
+ retryOptions: {now: clock.now, sleep: clock.sleep, random: () => 0},
+ })
+
+ assert.equal(data.length, 2)
+ assert.equal(calls, 2)
+ assert.equal(clock.t, 12000)
+ assert.deepEqual(rec.stages, {0: [0, 1], 1: [0, 1]})
+ assert.deepEqual(rec.failed, [])
+ assert.deepEqual(rec.finished, [0, 1])
+ assert.deepEqual(waits, ['throttled', 'end'])
+})
+
+test('runMatchBatch: a batch still throttled after the cap is marked throttled (-4), not failed', async () => {
+ const rec = recorder()
+ const clock = virtualClock()
+ const throttledRows = []
+
+ const data = await runMatchBatch({
+ rowIndexes: [5, 6], send: async () => { throw withRetryAfter(throttled(), 60) }, ...rec,
+ onRowThrottled: (index, info) => throttledRows.push([index, info.attempts > 1]),
+ retryOptions: {now: clock.now, sleep: clock.sleep, random: () => 0, maxWaitMs: 5 * 60 * 1000},
+ })
+
+ assert.deepEqual(data, [])
+ assert.deepEqual(rec.stages, {5: [0, -4], 6: [0, -4]})
+ assert.deepEqual(rec.failed, [])
+ assert.deepEqual(rec.finished, [])
+ assert.deepEqual(throttledRows, [[5, true], [6, true]])
+})
+
+test('runMatchBatch: Stop during a 429 wait puts the rows back to not run and sends nothing more', async () => {
+ const rec = recorder()
+ let calls = 0
+ let stopped = false
+ const send = async () => {
+ calls += 1
+ setTimeout(() => { stopped = true }, 20)
+ throw withRetryAfter(throttled(), 600)
+ }
+ const started = Date.now()
+
+ const data = await runMatchBatch({rowIndexes: [0, 1], send, ...rec, shouldStop: () => stopped, retryOptions: {pollMs: 5}})
+
+ assert.deepEqual(data, [])
+ assert.equal(calls, 1)
+ assert.ok(Date.now() - started < 1000)
+ assert.deepEqual(rec.stages, {0: [0, -1], 1: [0, -1]})
+ assert.deepEqual(rec.failed, [])
+})
+
+test('runMatchBatch: a 503 is still retried, then fails the rows with its status', async () => {
+ const rec = recorder()
+ let calls = 0
+ await runMatchBatch({rowIndexes: [0], send: async () => { calls += 1; throw httpError(503) }, ...rec, retryOptions: NO_WAIT})
+ assert.equal(calls, 3)
+ assert.deepEqual(rec.stages, {0: [0, -2]})
+ assert.equal(rec.failed[0].status, 503)
+})
+
+test('runMatchBatch: a 429 on one batch holds back another batch sharing the gate', async () => {
+ const clock = virtualClock()
+ const gate = createCapacityGate({now: clock.now})
+ const sentAt = {a: [], b: []}
+ const sendA = async () => {
+ sentAt.a.push(clock.t)
+ if(sentAt.a.length === 1) throw withRetryAfter(throttled(), 30)
+ return ok([])
+ }
+ const sendB = async () => { sentAt.b.push(clock.t); return ok([]) }
+ const common = {retryOptions: {gate, now: clock.now, sleep: clock.sleep, random: () => 0}}
+
+ await runMatchBatch({rowIndexes: [0], send: sendA, ...recorder(), ...common})
+ gate.pause(30000)
+ await runMatchBatch({rowIndexes: [1], send: sendB, ...recorder(), ...common})
+
+ assert.deepEqual(sentAt.a, [0, 30000])
+ assert.deepEqual(sentAt.b, [60000])
+})
+
+test('requestSingleMatch: a 429 waits and then returns the response', async () => {
+ const clock = virtualClock()
+ let calls = 0
+ const result = await requestSingleMatch(async () => {
+ calls += 1
+ if(calls === 1) throw withRetryAfter(throttled(), 5)
+ return ok([{ row: { __index: 3 }, results: [] }])
+ }, { retryOptions: {now: clock.now, sleep: clock.sleep, random: () => 0} })
+ assert.equal(result.ok, true)
+ assert.equal(calls, 2)
+ assert.equal(clock.t, 5000)
+})
+
+test('requestSingleMatch: still throttled after its cap, it reports throttled, not a failure', async () => {
+ const clock = virtualClock()
+ const result = await requestSingleMatch(async () => { throw withRetryAfter(throttled(), 60) }, { retryOptions: {now: clock.now, sleep: clock.sleep, random: () => 0} })
+ assert.equal(result.ok, false)
+ assert.equal(result.throttled, true)
+ assert.equal(result.previewLimit, false)
+ assert.equal(result.errorBody.status, 429)
+ // a person is waiting on the row panel: 5 minutes, not a run's 30
+ assert.ok(clock.t <= 5 * 60 * 1000 && clock.t >= 4 * 60 * 1000, `waited ${clock.t}`)
+})
+
+// ── the bulk scheduler ─────────────────────────────────────────────────────────
+
+test('runWithConcurrency: at most `concurrency` in flight, every item run once', async () => {
+ let inFlight = 0
+ let maxInFlight = 0
+ const ran = []
+ const done = await runWithConcurrency([1, 2, 3, 4, 5, 6, 7], {
+ concurrency: 3,
+ run: async item => {
+ inFlight += 1
+ maxInFlight = Math.max(maxInFlight, inFlight)
+ await new Promise(resolve => setTimeout(resolve, 3))
+ ran.push(item)
+ inFlight -= 1
+ },
+ })
+ assert.equal(done, true)
+ assert.equal(maxInFlight, 3)
+ assert.deepEqual(ran.sort(), [1, 2, 3, 4, 5, 6, 7])
+})
+
+test('runWithConcurrency: sends no new batch while the gate is paused, then resumes', async () => {
+ const gate = createCapacityGate()
+ const started = Date.now()
+ const dispatchedAt = {}
+ const done = await runWithConcurrency([1, 2, 3, 4], {
+ concurrency: 2,
+ gate,
+ pollMs: 5,
+ run: async item => {
+ dispatchedAt[item] = Date.now() - started
+ // the first batch is throttled: the server asks everyone to wait 80 ms
+ if(item === 1)
+ gate.pause(80)
+ await new Promise(resolve => setTimeout(resolve, 5))
+ },
+ })
+ assert.equal(done, true)
+ assert.ok(dispatchedAt[1] < 40)
+ for(const item of [2, 3, 4])
+ assert.ok(dispatchedAt[item] >= 75, `batch ${item} went out at ${dispatchedAt[item]} ms, during the pause`)
+})
+
+test('runWithConcurrency: Stop during a pause returns at once, and sends nothing more', async () => {
+ const gate = createCapacityGate()
+ let stopped = false
+ const ran = []
+ const started = Date.now()
+ setTimeout(() => { stopped = true }, 20)
+ const done = await runWithConcurrency([1, 2, 3], {
+ concurrency: 1,
+ gate,
+ pollMs: 5,
+ shouldAbort: () => stopped,
+ run: async item => { ran.push(item); gate.pause(60000) },
+ })
+ assert.equal(done, false)
+ assert.deepEqual(ran, [1])
+ assert.ok(Date.now() - started < 1000)
+})
+
+test('runWithConcurrency: shouldSkipRest drops the queue but lets the batches in flight finish', async () => {
+ let skip = false
+ const finished = []
+ const done = await runWithConcurrency([1, 2, 3, 4, 5], {
+ concurrency: 2,
+ shouldSkipRest: () => skip,
+ run: async item => {
+ if(item === 2) skip = true
+ await new Promise(resolve => setTimeout(resolve, 3))
+ finished.push(item)
+ },
+ })
+ assert.equal(done, true)
+ assert.deepEqual(finished.sort(), [1, 2])
+})
+
+test('runWithConcurrency: onSkipped gets the batches shouldSkipRest dropped, once', async () => {
+ let skip = false
+ const skipped = []
+ await runWithConcurrency([1, 2, 3, 4, 5], {
+ concurrency: 1,
+ shouldSkipRest: () => skip,
+ onSkipped: items => skipped.push(items),
+ run: async item => { if(item === 2) skip = true },
+ })
+ assert.deepEqual(skipped, [[3, 4, 5]])
+})
+
+test('runWithConcurrency: onSkipped isn\'t called when nothing was dropped', async () => {
+ const skipped = []
+ await runWithConcurrency([1, 2], {concurrency: 2, onSkipped: items => skipped.push(items), run: async () => {}})
+ assert.deepEqual(skipped, [])
+})
+
+// Codex review, pass 1
+test('runWithConcurrency: Stop returns even while a batch in flight never answers', async () => {
+ let stopped = false
+ setTimeout(() => { stopped = true }, 20)
+ const started = Date.now()
+ const done = await runWithConcurrency([1, 2, 3], {
+ concurrency: 1,
+ pollMs: 5,
+ shouldAbort: () => stopped,
+ run: () => new Promise(() => {}),
+ })
+ assert.equal(done, false)
+ assert.ok(Date.now() - started < 1000)
+})
+
+test('runWithConcurrency: a pause longer than maxHoldMs doesn\'t hold the queue (its batches end throttled at once instead)', async () => {
+ const gate = createCapacityGate()
+ gate.pause(60 * 60 * 1000)
+ const ran = []
+ const started = Date.now()
+ await runWithConcurrency([1, 2, 3], {concurrency: 1, gate, maxHoldMs: 30 * 60 * 1000, pollMs: 5, run: async item => { ran.push(item) }})
+ assert.deepEqual(ran, [1, 2, 3])
+ assert.ok(Date.now() - started < 1000)
+})
+
+// Codex review, pass 2: the scheduler's hold budget is cumulative.
+test('runWithConcurrency: a pause that keeps extending holds the queue for maxHoldMs in all, then lets it go', async () => {
+ const gate = createCapacityGate()
+ gate.pause(40)
+ // other requests keep pushing the pause out
+ const timer = setInterval(() => gate.pause(40), 10)
+ const started = Date.now()
+ const ran = []
+ await runWithConcurrency([1, 2], {concurrency: 1, gate, maxHoldMs: 150, pollMs: 5, run: async () => { ran.push(Date.now() - started) }})
+ clearInterval(timer)
+ assert.equal(ran.length, 2)
+ assert.ok(ran[0] >= 140 && ran[0] < 1000, `first batch went out at ${ran[0]} ms`)
+})
diff --git a/src/components/map-projects/__tests__/requestLimits.test.js b/src/components/map-projects/__tests__/requestLimits.test.js
index 78d7b85..795ee03 100644
--- a/src/components/map-projects/__tests__/requestLimits.test.js
+++ b/src/components/map-projects/__tests__/requestLimits.test.js
@@ -1,9 +1,12 @@
/**
* Auto Match request limits per algorithm (ocl_online#274).
*
- * Non-core users can't run more than 5 concurrent requests or batches of
- * more than 10 rows, whether they type it in the form or a project saved it
- * earlier. Core users and staff keep the form's full range.
+ * Non-core users can't run batches of more than 10 rows, whether they type it
+ * in the form or a project saved it earlier. Early-access users keep up to 5
+ * requests in flight; preview users 2 (ocl_issues#2849), as many previews at
+ * once could fill the server's matching capacity. A preview project has at
+ * most 25 rows (3 batches of 10), so 2 costs one more batch's time. Core users
+ * and staff keep the form's full range.
*
* Run with: npm test
*/
@@ -13,6 +16,10 @@ import assert from 'node:assert/strict'
import { hasFullRequestLimits, getRequestLimits, getRequestSettings, applyRequestSettings } from '../requestLimits.js'
+const FULL = getRequestLimits({ is_staff: true })
+const EARLY_ACCESS = getRequestLimits({ auth_groups: ['early_access'] })
+const PREVIEW = getRequestLimits({ auth_groups: ['preview'] })
+
test('hasFullRequestLimits: core users, staff and superusers', () => {
assert.equal(hasFullRequestLimits({ auth_groups: ['core_user'] }), true)
assert.equal(hasFullRequestLimits({ is_staff: true }), true)
@@ -27,52 +34,79 @@ test('hasFullRequestLimits: preview, early-access and signed-out users are cappe
assert.equal(hasFullRequestLimits(undefined), false)
})
-test('getRequestLimits: the form range for core users and staff, 10 rows × 5 requests for everyone else', () => {
- assert.deepEqual(getRequestLimits(true), { batchSize: 1000, concurrentRequests: 50 })
- assert.deepEqual(getRequestLimits(false), { batchSize: 10, concurrentRequests: 5 })
+test('getRequestLimits: the form range for core users and staff', () => {
+ assert.deepEqual(getRequestLimits({ auth_groups: ['core_user'] }), { batchSize: 1000, concurrentRequests: 50, capped: false })
+ assert.deepEqual(FULL, { batchSize: 1000, concurrentRequests: 50, capped: false })
+})
+
+test('getRequestLimits: early access, 10 rows × 5 requests', () => {
+ assert.deepEqual(EARLY_ACCESS, { batchSize: 10, concurrentRequests: 5, capped: true })
+})
+
+test('getRequestLimits: preview, 10 rows × 2 requests', () => {
+ assert.deepEqual(PREVIEW, { batchSize: 10, concurrentRequests: 2, capped: true })
+})
+
+test('getRequestLimits: a user in no tier, or signed out, gets the preview limits', () => {
+ assert.deepEqual(getRequestLimits({ auth_groups: [] }), PREVIEW)
+ assert.deepEqual(getRequestLimits({ auth_groups: ['mapper_ai_assistant'] }), PREVIEW)
+ assert.deepEqual(getRequestLimits(null), PREVIEW)
+})
+
+test('getRequestLimits: the highest tier wins', () => {
+ assert.deepEqual(getRequestLimits({ auth_groups: ['preview', 'early_access'] }), EARLY_ACCESS)
+ assert.deepEqual(getRequestLimits({ auth_groups: ['early_access', 'core_user'] }), FULL)
})
test('getRequestSettings: a capped user\'s saved project runs within the cap', () => {
// the ticket's case: batch 1 × 25 concurrent semantic requests
- assert.deepEqual(getRequestSettings({ batch_size: 1, concurrent_requests: 25 }, false), { batchSize: 1, concurrentRequests: 5 })
- assert.deepEqual(getRequestSettings({ batch_size: 1000, concurrent_requests: 50 }, false), { batchSize: 10, concurrentRequests: 5 })
+ assert.deepEqual(getRequestSettings({ batch_size: 1, concurrent_requests: 25 }, EARLY_ACCESS), { batchSize: 1, concurrentRequests: 5 })
+ assert.deepEqual(getRequestSettings({ batch_size: 1, concurrent_requests: 25 }, PREVIEW), { batchSize: 1, concurrentRequests: 2 })
+ assert.deepEqual(getRequestSettings({ batch_size: 1000, concurrent_requests: 50 }, EARLY_ACCESS), { batchSize: 10, concurrentRequests: 5 })
+ assert.deepEqual(getRequestSettings({ batch_size: 1000, concurrent_requests: 50 }, PREVIEW), { batchSize: 10, concurrentRequests: 2 })
// ocl-search's shipped default is 50 rows per batch
- assert.deepEqual(getRequestSettings({ batch_size: 50, concurrent_requests: 2 }, false), { batchSize: 10, concurrentRequests: 2 })
+ assert.deepEqual(getRequestSettings({ batch_size: 50, concurrent_requests: 2 }, EARLY_ACCESS), { batchSize: 10, concurrentRequests: 2 })
})
test('getRequestSettings: values within the cap are kept', () => {
- assert.deepEqual(getRequestSettings({ batch_size: 10, concurrent_requests: 2 }, false), { batchSize: 10, concurrentRequests: 2 })
- assert.deepEqual(getRequestSettings({ batch_size: 3, concurrent_requests: 5 }, false), { batchSize: 3, concurrentRequests: 5 })
+ assert.deepEqual(getRequestSettings({ batch_size: 10, concurrent_requests: 2 }, EARLY_ACCESS), { batchSize: 10, concurrentRequests: 2 })
+ assert.deepEqual(getRequestSettings({ batch_size: 3, concurrent_requests: 5 }, EARLY_ACCESS), { batchSize: 3, concurrentRequests: 5 })
+ assert.deepEqual(getRequestSettings({ batch_size: 3, concurrent_requests: 1 }, PREVIEW), { batchSize: 3, concurrentRequests: 1 })
})
test('getRequestSettings: core users and staff run what the project saved, even past the form\'s range', () => {
- assert.deepEqual(getRequestSettings({ batch_size: 1, concurrent_requests: 25 }, true), { batchSize: 1, concurrentRequests: 25 })
- assert.deepEqual(getRequestSettings({ batch_size: 50, concurrent_requests: 2 }, true), { batchSize: 50, concurrentRequests: 2 })
- assert.deepEqual(getRequestSettings({ batch_size: 2000, concurrent_requests: 75 }, true), { batchSize: 2000, concurrentRequests: 75 })
+ assert.deepEqual(getRequestSettings({ batch_size: 1, concurrent_requests: 25 }, FULL), { batchSize: 1, concurrentRequests: 25 })
+ assert.deepEqual(getRequestSettings({ batch_size: 50, concurrent_requests: 2 }, FULL), { batchSize: 50, concurrentRequests: 2 })
+ assert.deepEqual(getRequestSettings({ batch_size: 2000, concurrent_requests: 75 }, FULL), { batchSize: 2000, concurrentRequests: 75 })
})
test('getRequestSettings: missing or invalid values fall back to 10 rows × 1 request, as before', () => {
- for (const fullLimits of [true, false]) {
- assert.deepEqual(getRequestSettings({}, fullLimits), { batchSize: 10, concurrentRequests: 1 })
- assert.deepEqual(getRequestSettings(undefined, fullLimits), { batchSize: 10, concurrentRequests: 1 })
- assert.deepEqual(getRequestSettings({ batch_size: 0, concurrent_requests: 0 }, fullLimits), { batchSize: 10, concurrentRequests: 1 })
- assert.deepEqual(getRequestSettings({ batch_size: 'abc', concurrent_requests: -3 }, fullLimits), { batchSize: 10, concurrentRequests: 1 })
+ for (const limits of [FULL, EARLY_ACCESS, PREVIEW]) {
+ assert.deepEqual(getRequestSettings({}, limits), { batchSize: 10, concurrentRequests: 1 })
+ assert.deepEqual(getRequestSettings(undefined, limits), { batchSize: 10, concurrentRequests: 1 })
+ assert.deepEqual(getRequestSettings({ batch_size: 0, concurrent_requests: 0 }, limits), { batchSize: 10, concurrentRequests: 1 })
+ assert.deepEqual(getRequestSettings({ batch_size: 'abc', concurrent_requests: -3 }, limits), { batchSize: 10, concurrentRequests: 1 })
}
})
+test('getRequestSettings: without limits, the preview cap applies', () => {
+ assert.deepEqual(getRequestSettings({ batch_size: 50, concurrent_requests: 25 }), { batchSize: 10, concurrentRequests: 2 })
+})
+
test('getRequestSettings: numeric strings are read as numbers', () => {
- assert.deepEqual(getRequestSettings({ batch_size: '8', concurrent_requests: '9' }, false), { batchSize: 8, concurrentRequests: 5 })
+ assert.deepEqual(getRequestSettings({ batch_size: '8', concurrent_requests: '9' }, EARLY_ACCESS), { batchSize: 8, concurrentRequests: 5 })
})
test('applyRequestSettings: a capped user\'s algorithm carries the settings the run uses, for the run and its record', () => {
const algo = { id: 'ocl-semantic', type: 'ocl-semantic', batch_size: 1, concurrent_requests: 25, query_params: { semantic: true } }
- assert.deepEqual(applyRequestSettings(algo, false), { ...algo, batch_size: 1, concurrent_requests: 5 })
- assert.deepEqual(applyRequestSettings({ id: 'ocl-search' }, false), { id: 'ocl-search', batch_size: 10, concurrent_requests: 1 })
+ assert.deepEqual(applyRequestSettings(algo, EARLY_ACCESS), { ...algo, batch_size: 1, concurrent_requests: 5 })
+ assert.deepEqual(applyRequestSettings(algo, PREVIEW), { ...algo, batch_size: 1, concurrent_requests: 2 })
+ assert.deepEqual(applyRequestSettings({ id: 'ocl-search' }, EARLY_ACCESS), { id: 'ocl-search', batch_size: 10, concurrent_requests: 1 })
// the saved project is not changed
assert.equal(algo.concurrent_requests, 25)
})
test('applyRequestSettings: core users and staff get the algorithm back untouched', () => {
const algo = { id: 'ocl-semantic', batch_size: 2000, concurrent_requests: 75 }
- assert.equal(applyRequestSettings(algo, true), algo)
+ assert.equal(applyRequestSettings(algo, FULL), algo)
})
diff --git a/src/components/map-projects/__tests__/rowProgress.test.js b/src/components/map-projects/__tests__/rowProgress.test.js
new file mode 100644
index 0000000..a94a37c
--- /dev/null
+++ b/src/components/map-projects/__tests__/rowProgress.test.js
@@ -0,0 +1,73 @@
+/**
+ * The row panel's progress chip (Candidates.jsx).
+ *
+ * ocl_issues#2849: a row whose request is waiting out a busy server says so,
+ * and a row the server stayed too busy for says "retry", never "failed".
+ *
+ * Run with: npm test
+ */
+
+import test from 'node:test'
+import assert from 'node:assert/strict'
+
+import { getRowProgressLabel } from '../rowProgress.js'
+
+const ALGOS = [{id: 'ocl-semantic'}, {id: 'ocl-bridge'}]
+const t = key => key
+
+test('getRowProgressLabel: nothing for a row with no stages yet', () => {
+ assert.deepEqual(getRowProgressLabel(undefined, ALGOS, {t}), {label: false})
+})
+
+test('getRowProgressLabel: not started, running, waiting and partial, as before', () => {
+ assert.deepEqual(getRowProgressLabel({'ocl-semantic': -1, 'ocl-bridge': -1}, ALGOS, {t}), {label: 'Not started', status: 'idle'})
+ assert.deepEqual(getRowProgressLabel({'ocl-semantic': 0, 'ocl-bridge': -1}, ALGOS, {t}), {label: 'Running: ocl-semantic...', status: 'running'})
+ assert.deepEqual(getRowProgressLabel({'ocl-semantic': 1, 'ocl-bridge': -1}, ALGOS, {t}), {label: 'Waiting: ocl-bridge', status: 'waiting'})
+ assert.deepEqual(getRowProgressLabel({'ocl-semantic': 1, 'ocl-bridge': -2}, ALGOS, {t}), {label: 'Partially completed', status: 'partial'})
+})
+
+test('getRowProgressLabel: all done gives no label', () => {
+ assert.equal(getRowProgressLabel({'ocl-semantic': 1, 'ocl-bridge': 1}, ALGOS, {t}).label, undefined)
+})
+
+test('getRowProgressLabel: a row waiting for capacity says so, over "Running"', () => {
+ assert.deepEqual(
+ getRowProgressLabel({'ocl-semantic': 0, 'ocl-bridge': -1}, ALGOS, {t, capacityWait: true}),
+ {label: 'map_project.waiting_for_capacity', status: 'capacity_wait'},
+ )
+})
+
+test('getRowProgressLabel: a row waiting for capacity says so even once its algorithms are done (its rerank is waiting)', () => {
+ assert.equal(getRowProgressLabel({'ocl-semantic': 1, 'ocl-bridge': 1}, ALGOS, {t, capacityWait: true}).status, 'capacity_wait')
+})
+
+test('getRowProgressLabel: a throttled row asks for a retry, not "failed" or "partial"', () => {
+ assert.deepEqual(
+ getRowProgressLabel({'ocl-semantic': -4, 'ocl-bridge': 1}, ALGOS, {t}),
+ {label: 'map_project.row_throttled', status: 'throttled'},
+ )
+})
+
+test('getRowProgressLabel: while other algorithms still run, the row shows them, not throttled', () => {
+ assert.equal(getRowProgressLabel({'ocl-semantic': -4, 'ocl-bridge': 0}, ALGOS, {t}).status, 'running')
+ assert.equal(getRowProgressLabel({'ocl-semantic': -4, 'ocl-bridge': -1}, ALGOS, {t}).status, 'waiting')
+})
+
+test('getRowProgressLabel: an algorithm that can\'t run for this user (-3, n/a) counts as done', () => {
+ assert.equal(getRowProgressLabel({'ocl-semantic': 1, 'ocl-bridge': -3}, ALGOS, {t}).label, undefined)
+ assert.equal(getRowProgressLabel({'ocl-semantic': -2, 'ocl-bridge': -3}, ALGOS, {t}).status, 'partial')
+})
+
+// Codex review, pass 1: the rerank's own stage counts too.
+test('getRowProgressLabel: a throttled rerank asks for a retry once the algorithms are done', () => {
+ assert.deepEqual(
+ getRowProgressLabel({'ocl-semantic': 1, 'ocl-bridge': 1, rerank: -4}, ALGOS, {t}),
+ {label: 'map_project.row_throttled', status: 'throttled'},
+ )
+ assert.equal(getRowProgressLabel({'ocl-semantic': 1, 'ocl-bridge': -3, rerank: -4}, ALGOS, {t}).status, 'throttled')
+})
+
+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)
+})
diff --git a/src/components/map-projects/autoMatchRows.js b/src/components/map-projects/autoMatchRows.js
index 22df09c..77ab351 100644
--- a/src/components/map-projects/autoMatchRows.js
+++ b/src/components/map-projects/autoMatchRows.js
@@ -122,6 +122,16 @@ export const hasCurrentAnalysis = (analyses, fingerprint) => {
// the row's index, not its position in the run (ocl_online#258).
export const getScispacyRowResults = (responseData, rowIndex) => responseData?.[rowIndex] || []
+// Whether a failed ScispaCy call means the service is still starting: it
+// answers 503 {status: 'warming_up'} while its host boots, then a 502 until
+// the app is up. The row waits for it rather than failing.
+export const isScispacyWarmingUp = (err, seenWarmingUp = false) => {
+ const data = err?.response?.data
+ if(data?.status === 'warming_up')
+ return true
+ return Boolean(seenWarmingUp && err?.response?.status === 502)
+}
+
// The in-flight $lookups for a row's concepts. Rerank waits for them, and so
// does a bulk run's AI step, which can reach a row whose rerank was skipped
// (after a rerank quota stop), so the AI sees the looked-up concepts.
@@ -133,6 +143,11 @@ export const getPendingRowLookups = (rowState, inFlightLookups) =>
// How long a bulk run's AI step waits for a row's lookups before going ahead.
export const AI_LOOKUP_WAIT_MS = 15000
+// $rerank calls in flight at once, per tab, whatever fired them: the rerank
+// sweep, or the per-row reranks as batches finish (ocl_issues#2849). Each ties
+// up an API worker for seconds, and a finished 10-row batch used to fire ten.
+export const RERANK_MAX_IN_FLIGHT = 2
+
// How long rerank waits for a row's lookups before going ahead without them
// (ocl_online#283). A guard against a lookup that never settles, longer than
// the AI step's wait: a rerank that goes ahead early drops the concepts still
diff --git a/src/components/map-projects/autosave.js b/src/components/map-projects/autosave.js
index 30cac9a..1b32b6e 100644
--- a/src/components/map-projects/autosave.js
+++ b/src/components/map-projects/autosave.js
@@ -141,3 +141,40 @@ export const installUnloadGuard = () => {
event.returnValue = ''
})
}
+
+/**
+ * One send at a time; triggers that arrive during it coalesce into one more
+ * send (ocl_issues#2849). For the project logs POST, which sends the whole
+ * log: once it waits out a 429, an older POST could otherwise land after a
+ * newer one and overwrite it. send(getPayload) reads the latest payload
+ * triggered on this sender, on each attempt; keep one sender per project, so
+ * a retry for one project never sends another's. trigger(payload) resolves,
+ * never rejects, once nothing is left to send.
+ */
+export const createLatestSender = send => {
+ let chain = null
+ let pending = false
+ let latest
+ const getPayload = () => latest
+ const loop = async () => {
+ do {
+ pending = false
+ try {
+ await send(getPayload)
+ } catch (_) {
+ // Best effort: the next change sends the whole log again.
+ }
+ } while(pending)
+ chain = null
+ }
+ return {
+ trigger: payload => {
+ latest = payload
+ if(chain)
+ pending = true
+ else
+ chain = loop()
+ return chain
+ },
+ }
+}
diff --git a/src/components/map-projects/constants.jsx b/src/components/map-projects/constants.jsx
index fd8a711..a24146f 100644
--- a/src/components/map-projects/constants.jsx
+++ b/src/components/map-projects/constants.jsx
@@ -66,6 +66,8 @@ export const ES_BATCH_SIZE = 50
export const CANDIDATES_LIMIT = 30
export const ROW_STAGES = {
+ // the server stayed too busy: not run, retry (ocl_issues#2849)
+ '-4': 'throttled',
'-3': 'na',
'-2': 'failed',
'-1': 'not_started',
diff --git a/src/components/map-projects/matchBatch.js b/src/components/map-projects/matchBatch.js
index dd25732..1352f68 100644
--- a/src/components/map-projects/matchBatch.js
+++ b/src/components/map-projects/matchBatch.js
@@ -6,18 +6,22 @@
* candidates. The caller's `send` must reject instead (service.request does),
* so a failure is caught here, retried when it's transient, and reported per
* row.
+ *
+ * A 429 waits for the server's Retry-After and resumes (ocl_issues#2849): a
+ * busy server makes a run slower, not broken.
*/
import { isPreviewLimitError } from './previewLimits.js'
-import { isTransientNetworkError, retryWithBackoff } from '../../services/retry.js'
+import { isRetryableError, requestWithCapacityRetry, untilCancelled } from '../../services/capacity.js'
-// The gateway answers these while the API rolls out or a task restarts. A 500
-// comes from the API itself, and retrying it only adds load.
-const RETRYABLE_STATUSES = new Set([502, 503, 504])
+// How long the row panel waits out a 429 before showing the row as throttled.
+// A person is watching it, so less than a run's 30 minutes.
+export const INTERACTIVE_WAIT_CAP_MS = 5 * 60 * 1000
export const MATCH_RETRY_OPTIONS = {maxRetries: 2, baseDelayMs: 3000, backoffFactor: 4, jitterFactor: 0.25}
-export const isRetryableMatchError = err =>
- isTransientNetworkError(err) || RETRYABLE_STATUSES.has(err?.response?.status)
+// Network errors and 502/503/504. A 429 isn't an error to retry: it's waited
+// out, however many times it comes.
+export const isRetryableMatchError = isRetryableError
const getFailure = (err, attempts) => ({
error: err?.response?.data?.detail || err?.message || 'unknown',
@@ -27,13 +31,19 @@ const getFailure = (err, attempts) => ({
})
/**
- * Mark the batch's rows running, send it (retrying transient failures), then
- * mark each row done (1) or failed (-2). Never rejects.
+ * Mark the batch's rows running, send it, then mark each row done (1), failed
+ * (-2) or throttled (-4). Never rejects.
+ *
+ * A network error or a 502/503/504 is retried a bounded number of times, then
+ * fails the rows. A 429 waits for Retry-After, however often it comes, until
+ * the long cap; a batch still refused then is throttled (-4): "not run, retry",
+ * never failed. Waits on a gate shared with the run's other requests hold them
+ * all back together.
*
- * shouldStop is checked before every retry, including after the backoff
- * sleep. A batch stopped there (the run was cancelled, or another batch hit
- * the match quota) sends nothing more and puts its rows back to not run (-1),
- * like the rows the run never reached.
+ * shouldStop is checked before every send and during every wait. A batch
+ * stopped there (the run was cancelled, or another batch hit the match quota)
+ * sends nothing more and puts its rows back to not run (-1), like the rows the
+ * run never reached.
*
* @param {object} opts
* @param {number[]} opts.rowIndexes the batch's row __index values
@@ -41,57 +51,106 @@ const getFailure = (err, attempts) => ({
* @param {function} opts.setStage (rowIndex, stage) => void
* @param {function} opts.onRowFinished rowIndex => void
* @param {function} opts.onRowFailed (rowIndex, {error, status, attempts, previewLimit}) => void
+ * @param {function} [opts.onRowThrottled] (rowIndex, {attempts, waitedMs}) => void
* @param {function} [opts.onPreviewLimit] err => void, for a preview-limit 403 (never retried)
- * @param {function} [opts.shouldStop] () => boolean; true skips any further retry
- * @param {object} [opts.retryOptions] overrides MATCH_RETRY_OPTIONS
- * @returns {Promise