Skip to content

Commit 376a743

Browse files
authored
fix: eliminate Promise.race handler-stacking in batchPackageStream (#600)
* fix: eliminate Promise.race handler-stacking in batchPackageStream The prior implementation stored per-generator promises in a Map and called `Promise.race(running.values())` every loop iteration. Each race call re-attached fresh `.then` handlers to every still-pending promise in the pool, so long-surviving generators accumulated O(iterations) dead handler closures before finally settling. See nodejs/node#17469 and https://github.com/cefn/watchable/tree/main/packages/unpromise for the pattern. Switch to a single-waiter queue: each generator's `.then` delivers its step into a buffer, and the main loop awaits a fresh promiseWithResolvers each iteration. Handlers are one-shot — nothing to stack. All 565 tests pass. * docs(claude): add Promise.race handler-stacking rule Pairs with the batchPackageStream fix: documents the anti-pattern so future code does not reintroduce it. Concise bullet in the SHARED STANDARDS section alongside the existsSync rule. * docs(claude): drop duplicate Promise.race rule from TypeScript Patterns The rule already lives in SHARED STANDARDS (line 26). Having it twice — once general, once under sdk-specific TypeScript Patterns — adds maintenance drift without adding clarity. Keep the SHARED STANDARDS version. * style(sdk): alphabetize destructuring and prefer undefined in batchPackageStream Align the single-waiter queue with project rules: alphabetical destructuring/type property order and undefined (not null) for absent state.
1 parent 8fd7ae0 commit 376a743

2 files changed

Lines changed: 51 additions & 15 deletions

File tree

CLAUDE.md

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -16,6 +16,7 @@
1616
- 🚨 **NEVER use `npx`, `pnpm dlx`, or `yarn dlx`** — use `pnpm exec <package>` or `pnpm run <script>`. Add tools as pinned devDependencies first.
1717
- **minimumReleaseAge**: NEVER add packages to `minimumReleaseAgeExclude` in CI. Locally, ASK before adding — the age threshold is a security control.
1818
- File existence: ALWAYS `existsSync` from `node:fs`. NEVER `fs.access`, `fs.stat`-for-existence, or an async `fileExists` wrapper. Import form: `import { existsSync, promises as fs } from 'node:fs'`.
19+
- `Promise.race` / `Promise.any`: NEVER pass a long-lived promise (interrupt signal, pool member) into a race inside a loop. Each call re-attaches `.then` handlers to every arm; handlers accumulate on surviving promises until they settle. For concurrency limiters, use a single-waiter "slot available" signal (resolved by each task's `.then`) instead of re-racing `executing[]`. See nodejs/node#17469 and `@watchable/unpromise`. Race with two fresh arms (e.g. one-shot `withTimeout`) is safe.
1920

2021
---
2122

src/socket-sdk-class.ts

Lines changed: 50 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -867,10 +867,52 @@ export class SocketSdk {
867867
/* c8 ignore stop */
868868
const { components } = componentsObj
869869
const { length: componentsCount } = components
870-
const running = new Map<
871-
AsyncGenerator<BatchPackageFetchResultType>,
872-
Promise<GeneratorStep>
873-
>()
870+
// Tracks in-flight generators only for pool-size accounting.
871+
// Completed steps and errors flow through the single-waiter queue below,
872+
// not through per-generator promises re-raced each iteration — repeated
873+
// Promise.race() over the same pool accumulates unreleased .then
874+
// handlers on each still-pending arm until the pool drains.
875+
// See https://github.com/nodejs/node/issues/17469.
876+
const running = new Set<AsyncGenerator<BatchPackageFetchResultType>>()
877+
const completed: GeneratorStep[] = []
878+
let waiter:
879+
| {
880+
reject: (err: unknown) => void
881+
resolve: (step: GeneratorStep) => void
882+
}
883+
| undefined
884+
let pendingError: { err: unknown } | undefined
885+
const deliverStep = (step: GeneratorStep) => {
886+
if (waiter) {
887+
const w = waiter
888+
waiter = undefined
889+
w.resolve(step)
890+
} else {
891+
completed.push(step)
892+
}
893+
}
894+
const deliverError = (err: unknown) => {
895+
if (waiter) {
896+
const w = waiter
897+
waiter = undefined
898+
w.reject(err)
899+
} else if (!pendingError) {
900+
pendingError = { err }
901+
}
902+
}
903+
const takeStep = (): Promise<GeneratorStep> => {
904+
if (pendingError) {
905+
const { err } = pendingError
906+
pendingError = undefined
907+
return Promise.reject(err)
908+
}
909+
if (completed.length) {
910+
return Promise.resolve(completed.shift()!)
911+
}
912+
const { promise, reject, resolve } = promiseWithResolvers<GeneratorStep>()
913+
waiter = { reject, resolve }
914+
return promise
915+
}
874916
let index = 0
875917
const enqueueGen = () => {
876918
if (index >= componentsCount) {
@@ -888,17 +930,12 @@ export class SocketSdk {
888930
const continueGen = (
889931
generator: AsyncGenerator<BatchPackageFetchResultType>,
890932
) => {
891-
const {
892-
promise,
893-
reject: rejectFn,
894-
resolve: resolveFn,
895-
} = promiseWithResolvers<GeneratorStep>()
896-
running.set(generator, promise)
933+
running.add(generator)
897934
void generator
898935
.next()
899936
.then(
900-
iteratorResult => resolveFn({ generator, iteratorResult }),
901-
rejectFn,
937+
iteratorResult => deliverStep({ generator, iteratorResult }),
938+
deliverError,
902939
)
903940
}
904941
// Start initial batch of generators.
@@ -907,9 +944,7 @@ export class SocketSdk {
907944
}
908945
while (running.size > 0) {
909946
// eslint-disable-next-line no-await-in-loop
910-
const { generator, iteratorResult }: GeneratorStep = await Promise.race(
911-
running.values(),
912-
)
947+
const { generator, iteratorResult }: GeneratorStep = await takeStep()
913948
running.delete(generator)
914949
// Yield the value if one is given, even when done:true.
915950
if (iteratorResult.value) {

0 commit comments

Comments
 (0)