Skip to content

Commit bbed226

Browse files
committed
Add non-blocking mode
1 parent cf023cc commit bbed226

7 files changed

Lines changed: 579 additions & 58 deletions

File tree

README.md

Lines changed: 15 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -285,13 +285,14 @@ const consumer = new MySnsSqsConsumer(dependencies, {
285285
| `enabled` | `boolean` | - | Must be set to `true` to enable polling |
286286
| `pollingIntervalMs` | `number` | `5000` | Interval between availability checks (ms) |
287287
| `timeoutMs` | `number \| NO_TIMEOUT` | - (required) | Maximum wait time before throwing `StartupResourcePollingTimeoutError`. Use `NO_TIMEOUT` to poll indefinitely |
288+
| `throwOnTimeout` | `boolean` | `true` | When `true`, throws error on timeout. When `false`, reports error via errorReporter, resets timeout counter, and continues polling |
288289

289290
### Environment-Specific Configuration
290291

291292
```typescript
292293
import { NO_TIMEOUT } from '@message-queue-toolkit/core'
293294

294-
// Production - with timeout
295+
// Production - with timeout (throws error when timeout reached)
295296
{
296297
locatorConfig: {
297298
queueUrl: '...',
@@ -302,6 +303,19 @@ import { NO_TIMEOUT } from '@message-queue-toolkit/core'
302303
}
303304
}
304305

306+
// Production - report timeout but keep trying
307+
// Useful when you want visibility into prolonged unavailability without failing
308+
{
309+
locatorConfig: {
310+
queueUrl: '...',
311+
startupResourcePolling: {
312+
enabled: true,
313+
timeoutMs: 5 * 60 * 1000, // Report every 5 minutes
314+
throwOnTimeout: false, // Don't throw, just report and continue
315+
}
316+
}
317+
}
318+
305319
// Development/Staging - poll indefinitely
306320
{
307321
locatorConfig: {

packages/core/lib/types/MessageQueueTypes.ts

Lines changed: 6 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,8 @@
1-
import type { CommonLogger, TransactionObservabilityManager } from '@lokalise/node-core'
1+
import type {
2+
CommonLogger,
3+
ErrorReporter,
4+
TransactionObservabilityManager,
5+
} from '@lokalise/node-core'
26
import type { ZodSchema } from 'zod/v4'
37

48
import type { PublicHandlerSpy } from '../queues/HandlerSpy.ts'
@@ -38,6 +42,7 @@ export type { TransactionObservabilityManager }
3842

3943
export type ExtraParams = {
4044
logger?: CommonLogger
45+
errorReporter?: ErrorReporter
4146
}
4247

4348
export type SchemaMap<SupportedMessageTypes extends string> = Record<

packages/core/lib/types/queueOptionsTypes.ts

Lines changed: 22 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -213,6 +213,28 @@ export type StartupResourcePollingConfig = {
213213
* Default: 5000 (5 seconds)
214214
*/
215215
pollingIntervalMs?: number
216+
217+
/**
218+
* Whether to throw an error when timeout is reached.
219+
* - `true` (default): Throws `StartupResourcePollingTimeoutError` when timeout is reached
220+
* - `false`: Reports the error via errorReporter, resets the timeout counter, and continues polling
221+
*
222+
* Use `false` when you want to be notified about prolonged unavailability but don't want to fail.
223+
* Default: true
224+
*/
225+
throwOnTimeout?: boolean
226+
227+
/**
228+
* Whether to run polling in non-blocking mode.
229+
* - `false` (default): init() waits for the resource to become available before resolving
230+
* - `true`: If resource is not immediately available, init() resolves immediately and
231+
* polling continues in the background. When the resource becomes available,
232+
* the `onResourceAvailable` callback is invoked.
233+
*
234+
* Use `true` when you want the service to start quickly without waiting for dependencies.
235+
* Default: false
236+
*/
237+
nonBlocking?: boolean
216238
}
217239

218240
/**

packages/core/lib/utils/startupResourcePollingUtils.spec.ts

Lines changed: 228 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -156,6 +156,234 @@ describe('startupResourcePollingUtils', () => {
156156
expect(result).toBe('test-result')
157157
expect(checkFn).toHaveBeenCalledTimes(5)
158158
})
159+
160+
it('throws by default when throwOnTimeout is not specified', async () => {
161+
const checkFn = vi.fn().mockResolvedValue({ isAvailable: false })
162+
163+
await expect(
164+
waitForResource({
165+
config: { enabled: true, pollingIntervalMs: 10, timeoutMs: 50 },
166+
checkFn,
167+
resourceName: 'test-resource',
168+
}),
169+
).rejects.toThrow(StartupResourcePollingTimeoutError)
170+
})
171+
172+
it('throws when throwOnTimeout is true', async () => {
173+
const checkFn = vi.fn().mockResolvedValue({ isAvailable: false })
174+
175+
await expect(
176+
waitForResource({
177+
config: { enabled: true, pollingIntervalMs: 10, timeoutMs: 50, throwOnTimeout: true },
178+
checkFn,
179+
resourceName: 'test-resource',
180+
}),
181+
).rejects.toThrow(StartupResourcePollingTimeoutError)
182+
})
183+
184+
it('reports error and continues polling when throwOnTimeout is false', async () => {
185+
let callCount = 0
186+
const checkFn = vi.fn().mockImplementation(() => {
187+
callCount++
188+
// Become available after what would be 2 timeout cycles
189+
if (callCount < 8) {
190+
return Promise.resolve({ isAvailable: false })
191+
}
192+
return Promise.resolve({ isAvailable: true, result: 'test-result' })
193+
})
194+
195+
const errorReporter = {
196+
report: vi.fn(),
197+
}
198+
199+
const result = await waitForResource({
200+
config: {
201+
enabled: true,
202+
pollingIntervalMs: 10,
203+
timeoutMs: 35, // Will timeout around every 3-4 calls
204+
throwOnTimeout: false,
205+
},
206+
checkFn,
207+
resourceName: 'test-resource',
208+
errorReporter,
209+
})
210+
211+
expect(result).toBe('test-result')
212+
// Should have reported at least one timeout error
213+
expect(errorReporter.report).toHaveBeenCalled()
214+
expect(errorReporter.report).toHaveBeenCalledWith(
215+
expect.objectContaining({
216+
error: expect.any(StartupResourcePollingTimeoutError),
217+
context: expect.objectContaining({
218+
resourceName: 'test-resource',
219+
timeoutMs: 35,
220+
}),
221+
}),
222+
)
223+
})
224+
225+
it('logs warning when throwOnTimeout is false and timeout is reached', async () => {
226+
let callCount = 0
227+
const checkFn = vi.fn().mockImplementation(() => {
228+
callCount++
229+
if (callCount < 5) {
230+
return Promise.resolve({ isAvailable: false })
231+
}
232+
return Promise.resolve({ isAvailable: true, result: 'test-result' })
233+
})
234+
235+
const logger = {
236+
info: vi.fn(),
237+
debug: vi.fn(),
238+
error: vi.fn(),
239+
warn: vi.fn(),
240+
}
241+
242+
await waitForResource({
243+
config: {
244+
enabled: true,
245+
pollingIntervalMs: 10,
246+
timeoutMs: 25,
247+
throwOnTimeout: false,
248+
},
249+
checkFn,
250+
resourceName: 'test-resource',
251+
// @ts-expect-error - partial logger for testing
252+
logger,
253+
})
254+
255+
expect(logger.warn).toHaveBeenCalledWith(
256+
expect.objectContaining({
257+
message: expect.stringContaining('resetting timeout counter'),
258+
resourceName: 'test-resource',
259+
}),
260+
)
261+
})
262+
263+
it('tracks timeout count when throwOnTimeout is false', async () => {
264+
let callCount = 0
265+
const checkFn = vi.fn().mockImplementation(() => {
266+
callCount++
267+
// Become available after multiple timeout cycles
268+
if (callCount < 12) {
269+
return Promise.resolve({ isAvailable: false })
270+
}
271+
return Promise.resolve({ isAvailable: true, result: 'test-result' })
272+
})
273+
274+
const errorReporter = {
275+
report: vi.fn(),
276+
}
277+
278+
await waitForResource({
279+
config: {
280+
enabled: true,
281+
pollingIntervalMs: 5,
282+
timeoutMs: 20, // Will timeout every ~4 calls
283+
throwOnTimeout: false,
284+
},
285+
checkFn,
286+
resourceName: 'test-resource',
287+
errorReporter,
288+
})
289+
290+
// Should have reported multiple timeout errors with increasing timeoutCount
291+
expect(errorReporter.report.mock.calls.length).toBeGreaterThan(1)
292+
293+
// Verify timeoutCount increments
294+
const firstCall = errorReporter.report.mock.calls[0]?.[0] as {
295+
context: { timeoutCount: number }
296+
}
297+
const secondCall = errorReporter.report.mock.calls[1]?.[0] as {
298+
context: { timeoutCount: number }
299+
}
300+
expect(firstCall.context.timeoutCount).toBe(1)
301+
expect(secondCall.context.timeoutCount).toBe(2)
302+
})
303+
304+
describe('nonBlocking mode', () => {
305+
it('returns result immediately when resource is available on first check', async () => {
306+
const checkFn = vi.fn().mockResolvedValue({ isAvailable: true, result: 'test-result' })
307+
const onResourceAvailable = vi.fn()
308+
309+
const result = await waitForResource({
310+
config: { enabled: true, pollingIntervalMs: 100, timeoutMs: 5000, nonBlocking: true },
311+
checkFn,
312+
resourceName: 'test-resource',
313+
onResourceAvailable,
314+
})
315+
316+
expect(result).toBe('test-result')
317+
expect(checkFn).toHaveBeenCalledTimes(1)
318+
expect(onResourceAvailable).not.toHaveBeenCalled()
319+
})
320+
321+
it('returns undefined immediately when resource is not available and starts background polling', async () => {
322+
const checkFn = vi.fn().mockResolvedValue({ isAvailable: false })
323+
const onResourceAvailable = vi.fn()
324+
325+
const result = await waitForResource({
326+
config: { enabled: true, pollingIntervalMs: 100, timeoutMs: 5000, nonBlocking: true },
327+
checkFn,
328+
resourceName: 'test-resource',
329+
onResourceAvailable,
330+
})
331+
332+
expect(result).toBeUndefined()
333+
expect(checkFn).toHaveBeenCalledTimes(1)
334+
})
335+
336+
it('calls onResourceAvailable when resource becomes available in background', async () => {
337+
let callCount = 0
338+
const checkFn = vi.fn().mockImplementation(() => {
339+
callCount++
340+
if (callCount < 3) {
341+
return Promise.resolve({ isAvailable: false })
342+
}
343+
return Promise.resolve({ isAvailable: true, result: 'test-result' })
344+
})
345+
const onResourceAvailable = vi.fn()
346+
347+
const result = await waitForResource({
348+
config: { enabled: true, pollingIntervalMs: 10, timeoutMs: 5000, nonBlocking: true },
349+
checkFn,
350+
resourceName: 'test-resource',
351+
onResourceAvailable,
352+
})
353+
354+
expect(result).toBeUndefined()
355+
356+
// Wait for background polling to complete
357+
await new Promise((resolve) => globalThis.setTimeout(resolve, 100))
358+
359+
expect(onResourceAvailable).toHaveBeenCalledWith('test-result')
360+
})
361+
362+
it('logs info message when starting background polling', async () => {
363+
const checkFn = vi.fn().mockResolvedValue({ isAvailable: false })
364+
const logger = {
365+
info: vi.fn(),
366+
debug: vi.fn(),
367+
error: vi.fn(),
368+
warn: vi.fn(),
369+
}
370+
371+
await waitForResource({
372+
config: { enabled: true, pollingIntervalMs: 100, timeoutMs: 5000, nonBlocking: true },
373+
checkFn,
374+
resourceName: 'test-resource',
375+
// @ts-expect-error - partial logger for testing
376+
logger,
377+
})
378+
379+
expect(logger.info).toHaveBeenCalledWith(
380+
expect.objectContaining({
381+
message: expect.stringContaining('starting background polling'),
382+
resourceName: 'test-resource',
383+
}),
384+
)
385+
})
386+
})
159387
})
160388

161389
describe('StartupResourcePollingTimeoutError', () => {

0 commit comments

Comments
 (0)