Skip to content

Commit c5f331c

Browse files
authored
fix: rate limiting api call (CM-1138) (#4131)
Signed-off-by: Umberto Sgueglia <usgueglia@contractor.linuxfoundation.org>
1 parent b1c679c commit c5f331c

3 files changed

Lines changed: 119 additions & 42 deletions

File tree

services/apps/automatic_projects_discovery_worker/src/activities/activities.ts

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,3 +1,4 @@
1+
import { Context } from '@temporalio/activity'
12
import { parse } from 'csv-parse'
23

34
import { bulkUpsertProjectCatalog } from '@crowd/data-access-layer'
@@ -102,6 +103,7 @@ export async function processDataset(
102103
totalProcessed += batch.length
103104
batch = []
104105

106+
Context.current().heartbeat({ totalProcessed, batchNumber })
105107
log.info({ totalProcessed, batchNumber, datasetId: dataset.id }, 'Batch upserted.')
106108
}
107109
}

services/apps/automatic_projects_discovery_worker/src/sources/lf-criticality-score/source.ts

Lines changed: 114 additions & 40 deletions
Original file line numberDiff line numberDiff line change
@@ -2,6 +2,7 @@ import http from 'http'
22
import https from 'https'
33
import { Readable } from 'stream'
44

5+
import { timeout } from '@crowd/common'
56
import { getServiceLogger } from '@crowd/logging'
67

78
import { IDatasetDescriptor, IDiscoverySource, IDiscoverySourceRow } from '../types'
@@ -12,6 +13,31 @@ const DEFAULT_API_HOST = 'lf-criticality-score-api.example.com'
1213
const DEFAULT_API_PORT = 443
1314
const PAGE_SIZE = 100
1415

16+
function parseEnvInt(
17+
value: string | undefined,
18+
defaultValue: number,
19+
min: number,
20+
max: number,
21+
): number {
22+
const parsed = parseInt(value ?? '', 10)
23+
return Number.isFinite(parsed) && parsed >= min && parsed <= max ? parsed : defaultValue
24+
}
25+
26+
// Requests per second sent to the LF Criticality Score API (throttle between pages).
27+
const REQUESTS_PER_SECOND = parseEnvInt(
28+
process.env.LF_CRITICALITY_SCORE_REQUESTS_PER_SECOND,
29+
5,
30+
1,
31+
100,
32+
)
33+
// Max per-page attempts (initial + retries) on 429 or 5xx before giving up.
34+
const MAX_ATTEMPTS = parseEnvInt(
35+
process.env.LF_CRITICALITY_SCORE_MAX_ATTEMPTS ?? process.env.LF_CRITICALITY_SCORE_MAX_RETRIES,
36+
7,
37+
1,
38+
20,
39+
)
40+
1541
interface LfApiResponse {
1642
page: number
1743
pageSize: number
@@ -46,6 +72,37 @@ function getApiBaseUrl(): string {
4672
return `${scheme}://${host}:${port}`
4773
}
4874

75+
interface HttpGetResult {
76+
statusCode: number
77+
retryAfterMs: number | null
78+
body: string
79+
}
80+
81+
function parseRetryAfterMs(header: string | string[] | undefined): number | null {
82+
const raw = Array.isArray(header) ? header[0] : header
83+
if (!raw) return null
84+
const secs = parseFloat(raw.trim())
85+
return Number.isFinite(secs) && secs > 0 ? secs * 1000 : null
86+
}
87+
88+
function httpGet(url: string): Promise<HttpGetResult> {
89+
return new Promise((resolve, reject) => {
90+
const client = url.startsWith('https://') ? https : http
91+
const req = client.get(url, (res) => {
92+
const statusCode = res.statusCode ?? 0
93+
const retryAfterMs = parseRetryAfterMs(res.headers['retry-after'])
94+
const chunks: Uint8Array[] = []
95+
res.on('data', (chunk: Uint8Array) => chunks.push(chunk))
96+
res.on('end', () =>
97+
resolve({ statusCode, retryAfterMs, body: Buffer.concat(chunks).toString('utf8') }),
98+
)
99+
res.on('error', reject)
100+
})
101+
req.on('error', reject)
102+
req.end()
103+
})
104+
}
105+
49106
async function fetchPage(
50107
baseUrl: string,
51108
page: number,
@@ -55,31 +112,50 @@ async function fetchPage(
55112
if (scoredAfter) params.set('scoredAfter', scoredAfter)
56113
const url = `${baseUrl}/projects?${params.toString()}`
57114

58-
return new Promise((resolve, reject) => {
59-
const client = url.startsWith('https://') ? https : http
115+
for (let attempt = 0; attempt < MAX_ATTEMPTS; attempt++) {
116+
let result: HttpGetResult | null = null
60117

61-
const req = client.get(url, (res) => {
62-
if (res.statusCode !== 200) {
63-
reject(new Error(`LF Criticality Score API returned status ${res.statusCode} for ${url}`))
64-
res.resume()
65-
return
118+
try {
119+
result = await httpGet(url)
120+
} catch (networkErr) {
121+
if (attempt === MAX_ATTEMPTS - 1) {
122+
throw new Error(`LF Criticality Score API network error for ${url}: ${networkErr}`)
66123
}
124+
const delayMs = Math.min(Math.pow(2, attempt) * 1000, 60_000)
125+
log.warn(
126+
{ page, attempt: attempt + 1, maxAttempts: MAX_ATTEMPTS, delayMs, err: String(networkErr) },
127+
'LF Criticality Score: network error, retrying...',
128+
)
129+
await timeout(delayMs)
130+
continue
131+
}
67132

68-
const chunks: Uint8Array[] = []
69-
res.on('data', (chunk: Uint8Array) => chunks.push(chunk))
70-
res.on('end', () => {
71-
try {
72-
resolve(JSON.parse(Buffer.concat(chunks).toString('utf8')) as LfApiResponse)
73-
} catch (err) {
74-
reject(new Error(`Failed to parse LF Criticality Score API response: ${err}`))
75-
}
76-
})
77-
res.on('error', reject)
78-
})
133+
const { statusCode, retryAfterMs, body } = result
79134

80-
req.on('error', reject)
81-
req.end()
82-
})
135+
if (statusCode === 200) {
136+
try {
137+
return JSON.parse(body) as LfApiResponse
138+
} catch (err) {
139+
throw new Error(`Failed to parse LF Criticality Score API response: ${err}`)
140+
}
141+
}
142+
143+
const isRetryable = statusCode === 429 || statusCode >= 500
144+
if (!isRetryable || attempt === MAX_ATTEMPTS - 1) {
145+
throw new Error(`LF Criticality Score API returned status ${statusCode} for ${url}`)
146+
}
147+
148+
const delayMs = retryAfterMs ?? Math.min(Math.pow(2, attempt) * 1000, 60_000)
149+
150+
log.warn(
151+
{ page, attempt: attempt + 1, maxAttempts: MAX_ATTEMPTS, statusCode, delayMs },
152+
'LF Criticality Score: rate limited or server error, retrying...',
153+
)
154+
await timeout(delayMs)
155+
}
156+
157+
// Unreachable, but satisfies TypeScript.
158+
throw new Error(`LF Criticality Score API failed for ${url} after ${MAX_ATTEMPTS} attempts`)
83159
}
84160

85161
export class LfCriticalityScoreSource implements IDiscoverySource {
@@ -113,29 +189,29 @@ export class LfCriticalityScoreSource implements IDiscoverySource {
113189
'LF Criticality Score: starting stream fetch.',
114190
)
115191

192+
const throttleIntervalMs = Math.round(1000 / REQUESTS_PER_SECOND)
193+
116194
async function* pages() {
117-
let page = 1
118-
let totalPages = 1
195+
const firstPage = await fetchPage(baseUrl, 1, scoredAfter)
196+
const { totalPages } = firstPage
197+
198+
log.info(
199+
{ datasetId: dataset.id, total: firstPage.total, totalPages, pageSize: firstPage.pageSize },
200+
'LF Criticality Score: first page received — total records available.',
201+
)
202+
203+
for (const row of firstPage.data) {
204+
yield row
205+
}
206+
207+
for (let page = 2; page <= totalPages; page++) {
208+
await timeout(throttleIntervalMs)
119209

120-
do {
121210
log.info(
122211
{ datasetId: dataset.id, page, totalPages },
123212
'LF Criticality Score: fetching page...',
124213
)
125214
const response = await fetchPage(baseUrl, page, scoredAfter)
126-
totalPages = response.totalPages
127-
128-
if (page === 1) {
129-
log.info(
130-
{
131-
datasetId: dataset.id,
132-
total: response.total,
133-
totalPages,
134-
pageSize: response.pageSize,
135-
},
136-
'LF Criticality Score: first page received — total records available.',
137-
)
138-
}
139215

140216
for (const row of response.data) {
141217
yield row
@@ -145,9 +221,7 @@ export class LfCriticalityScoreSource implements IDiscoverySource {
145221
{ datasetId: dataset.id, page, totalPages, rowsInPage: response.data.length },
146222
'LF Criticality Score: page fetched.',
147223
)
148-
149-
page++
150-
} while (page <= totalPages)
224+
}
151225

152226
log.info({ datasetId: dataset.id, totalPages }, 'LF Criticality Score: all pages fetched.')
153227
}

services/apps/automatic_projects_discovery_worker/src/workflows/discoverProjects.ts

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -7,9 +7,10 @@ const listActivities = proxyActivities<typeof activities>({
77
retry: { maximumAttempts: 3 },
88
})
99

10-
// processDataset is long-running (10-20 min for ~119MB / ~750K rows).
10+
// processDataset is long-running: ~119MB / ~750K rows + inter-page throttle can exceed 60 min.
1111
const processActivities = proxyActivities<typeof activities>({
12-
startToCloseTimeout: '30 minutes',
12+
startToCloseTimeout: '90 minutes',
13+
heartbeatTimeout: '5 minutes',
1314
retry: { maximumAttempts: 3 },
1415
})
1516

0 commit comments

Comments
 (0)