Skip to content

Commit 0235cb9

Browse files
authored
Merge pull request #961 from d-zero-dev/feat/dealer-runtime-limit-and-footer
feat(dealer): support runtime concurrency changes and injected Lanes
2 parents 245a487 + f649613 commit 0235cb9

9 files changed

Lines changed: 478 additions & 38 deletions

File tree

‎packages/@d-zero/dealer/README.md‎

Lines changed: 24 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -39,6 +39,30 @@ await deal(items, setup, { limit: 10, signal: controller.signal });
3939

4040
abort 時の挙動: **新規ワーカー起動を停止、実行中ワーカーは完了まで待機、`push`/`unshift` は無視**。詳細は `src/deal.ts` / `src/dealer.ts` の JSDoc。
4141

42+
### 実行中の並列数変更・外部 `Lanes` の再利用
43+
44+
```ts
45+
import { deal, Lanes } from '@d-zero/dealer';
46+
47+
const lanes = new Lanes({ stream: process.stderr });
48+
let controller: DealController | undefined;
49+
50+
const run = deal(items, setup, {
51+
limit: 10,
52+
lanes, // 呼び出し元が生成した Lanes を使い回す(deal() は生成も破棄もしない)
53+
onStart: (c) => {
54+
controller = c;
55+
},
56+
});
57+
58+
// deal() 実行中(run が解決する前)に、別イベントから並列数を変更する
59+
onExternalCommand((newLimit) => controller?.setLimit(newLimit));
60+
61+
await run;
62+
```
63+
64+
`lanes` を渡すと `deal()` は自前で `Lanes` を作らず、渡されたインスタンスの生成・破棄は呼び出し元の責任になる。`onStart` は `dealer.play()` 直前に一度呼ばれ、`controller.setLimit(n)` を **`await deal(...)` が解決する前に** 呼ぶことで実行中に並列数を変更できる(`Lanes#footer(text)` と組み合わせれば、外部からの入力受付 UI を並列レーンの下に固定表示できる)。
65+
4266
## Sequential Pipeline(`TaskList`)
4367

4468
```ts

‎packages/@d-zero/dealer/src/deal.spec.ts‎

Lines changed: 71 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,9 @@
1+
import type { DealController } from './types.js';
2+
13
import { describe, test, expect, vi } from 'vitest';
24

35
import { deal } from './deal.js';
6+
import { Lanes } from './lanes.js';
47

58
/**
69
*
@@ -152,4 +155,72 @@ describe('deal', () => {
152155

153156
stdoutWriteSpy.mockRestore();
154157
});
158+
159+
test('reuses an injected Lanes instead of creating its own, and does not dispose it', async () => {
160+
const resizeBefore = process.stdout.listenerCount('resize');
161+
const lanes = new Lanes({ verbose: true });
162+
// injected の場合、deal() が独自の Lanes を作らないので resize リスナーは増えない
163+
expect(process.stdout.listenerCount('resize')).toBe(resizeBefore + 1);
164+
165+
await deal(
166+
createItems(2),
167+
(_process, update, index) => {
168+
return () => {
169+
update(`item ${index}`);
170+
};
171+
},
172+
{ limit: 10, lanes },
173+
);
174+
175+
// deal() 完了後も呼び出し元の Lanes は破棄されず生きている
176+
expect(process.stdout.listenerCount('resize')).toBe(resizeBefore + 1);
177+
lanes[Symbol.dispose]();
178+
expect(process.stdout.listenerCount('resize')).toBe(resizeBefore);
179+
});
180+
181+
test('onStart receives a controller before play(), and setLimit affects the header limit', async () => {
182+
const limits: number[] = [];
183+
let controller: DealController | undefined;
184+
const { promise: firstStarted, resolve: resolveFirstStarted } =
185+
Promise.withResolvers<void>();
186+
const { promise: canFinishFirst, resolve: resolveCanFinishFirst } =
187+
Promise.withResolvers<void>();
188+
let firstCall = true;
189+
190+
const run = deal(
191+
createItems(3),
192+
() => {
193+
return async () => {
194+
if (firstCall) {
195+
firstCall = false;
196+
resolveFirstStarted();
197+
await canFinishFirst;
198+
}
199+
};
200+
},
201+
{
202+
limit: 1,
203+
verbose: true,
204+
onStart: (c) => {
205+
controller = c;
206+
},
207+
header: (_progress, _done, _total, limit) => {
208+
limits.push(limit);
209+
return `limit: ${limit}`;
210+
},
211+
},
212+
);
213+
214+
await firstStarted;
215+
expect(controller).toBeDefined();
216+
expect(controller?.limit).toBe(1);
217+
218+
controller?.setLimit(3);
219+
expect(controller?.limit).toBe(3);
220+
resolveCanFinishFirst();
221+
222+
await run;
223+
224+
expect(limits).toContain(3);
225+
});
155226
});

‎packages/@d-zero/dealer/src/deal.ts‎

Lines changed: 25 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,6 @@
11
import type { DealerOptions } from './dealer.js';
22
import type { LanesOptions } from './lanes.js';
3+
import type { DealController } from './types.js';
34
import type { DelayOptions } from '@d-zero/shared/delay';
45

56
import { delay } from '@d-zero/shared/delay';
@@ -18,6 +19,19 @@ export type DealOptions<T = unknown> = DealerOptions<T> &
1819
readonly header?: DealHeader;
1920
readonly debug?: boolean;
2021
readonly interval?: number | DelayOptions;
22+
/**
23+
* 呼び出し元が既に持っている `Lanes` インスタンスを使い回す。
24+
* 指定した場合、`deal()` はこの `Lanes` を生成も破棄もしない —
25+
* 呼び出し元が生成・破棄のライフサイクルを管理する。
26+
* 指定時は他の {@link LanesOptions}(`stream`/`verbose`/`fps` 等)は
27+
* 無視される(渡された `Lanes` 自身の設定が使われるため)。
28+
*/
29+
readonly lanes?: Lanes;
30+
/**
31+
* `dealer.play()` の直前に一度だけ呼ばれ、実行中に同時実行数を
32+
* 操作できる {@link DealController} を渡す。
33+
*/
34+
readonly onStart?: (controller: DealController) => void;
2135
};
2236

2337
/**
@@ -98,7 +112,11 @@ export async function deal<T extends WeakKey>(
98112
const dealer = new Dealer(items, options);
99113
// `using` により、setup() が例外を投げてもスコープ脱出時に必ず
100114
// lanes(内部の Display)のタイマー・resize リスナー・SIGINT ハンドラが解放される。
101-
using lanes = new Lanes(options);
115+
// `options.lanes` が渡された場合は呼び出し元が生成・破棄を管理するため、
116+
// ここでは新規生成も dispose もしない(`using` は null/undefined を
117+
// 許容し、その場合 dispose を呼ばない)。
118+
using ownedLanes = options?.lanes ? undefined : new Lanes(options);
119+
const lanes = options?.lanes ?? ownedLanes!;
102120

103121
if (options?.header) {
104122
dealer.progress((progress, done, total, limit) => {
@@ -137,6 +155,12 @@ export async function deal<T extends WeakKey>(
137155
// 実行される必要がある)ため、ここは `await` で完了を待ってからスコープを抜ける。
138156
const { promise, resolve } = Promise.withResolvers<void>();
139157
dealer.finish(resolve);
158+
options?.onStart?.({
159+
get limit() {
160+
return dealer.limit;
161+
},
162+
setLimit: (limit) => dealer.setLimit(limit),
163+
});
140164
dealer.play();
141165
await promise;
142166
}

‎packages/@d-zero/dealer/src/dealer.spec.ts‎

Lines changed: 140 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -604,4 +604,144 @@ describe('Dealer', () => {
604604

605605
await runDealer(dealer);
606606
});
607+
608+
describe('setLimit', () => {
609+
test('increasing the limit immediately fills newly available slots', async () => {
610+
const items = createItems(4);
611+
let maxConcurrent = 0;
612+
let currentConcurrent = 0;
613+
const dealer = new Dealer(items, { limit: 1 });
614+
const { promise: firstStarted, resolve: resolveFirstStarted } =
615+
Promise.withResolvers<void>();
616+
const { promise: canFinishFirst, resolve: resolveCanFinishFirst } =
617+
Promise.withResolvers<void>();
618+
let firstCall = true;
619+
620+
await dealer.setup(() => {
621+
return Promise.resolve(async () => {
622+
currentConcurrent++;
623+
maxConcurrent = Math.max(maxConcurrent, currentConcurrent);
624+
if (firstCall) {
625+
firstCall = false;
626+
resolveFirstStarted();
627+
await canFinishFirst;
628+
}
629+
currentConcurrent--;
630+
});
631+
});
632+
633+
const done = runDealer(dealer);
634+
await firstStarted;
635+
// limit 1 のあいだは1件しか動いていないはず
636+
expect(maxConcurrent).toBe(1);
637+
638+
dealer.setLimit(4);
639+
resolveCanFinishFirst();
640+
await done;
641+
642+
expect(maxConcurrent).toBeGreaterThan(1);
643+
expect(dealer.limit).toBe(4);
644+
});
645+
646+
test('decreasing the limit lets already-started workers finish, but throttles concurrency for items dispatched afterward', async () => {
647+
// limit 3・5件: 最初の3件(index 0,1,2)が同時ディスパッチされる。
648+
// それらが実行中のうちに limit を 1 へ落とし、(a) 実行中の3件は
649+
// 中断されず全件完了すること、(b) 減少後に新規ディスパッチされる
650+
// 残り2件(index 3,4)は同時に1件までしか動かないことを検証する。
651+
const items = createItems(5);
652+
const dealer = new Dealer(items, { limit: 3 });
653+
const startedFirstBatch: number[] = [];
654+
const { promise: firstBatchStarted, resolve: resolveFirstBatchStarted } =
655+
Promise.withResolvers<void>();
656+
const { promise: canFinishFirstBatch, resolve: resolveCanFinishFirstBatch } =
657+
Promise.withResolvers<void>();
658+
let decreased = false;
659+
let concurrentAfterDecrease = 0;
660+
let maxConcurrentAfterDecrease = 0;
661+
662+
await dealer.setup((_item, index) => {
663+
return Promise.resolve(async () => {
664+
if (!decreased) {
665+
startedFirstBatch.push(index);
666+
if (startedFirstBatch.length === 3) {
667+
resolveFirstBatchStarted();
668+
}
669+
await canFinishFirstBatch;
670+
return;
671+
}
672+
concurrentAfterDecrease++;
673+
maxConcurrentAfterDecrease = Math.max(
674+
maxConcurrentAfterDecrease,
675+
concurrentAfterDecrease,
676+
);
677+
await new Promise((r) => setTimeout(r, 5));
678+
concurrentAfterDecrease--;
679+
});
680+
});
681+
682+
const done = runDealer(dealer);
683+
await firstBatchStarted;
684+
expect(startedFirstBatch).toHaveLength(3);
685+
686+
// setLimit はテストの非同期フロー(=ワーカー自身の同期区間の外)から
687+
// 呼ぶ。これは実運用(外部入力による並列数変更)と同じ呼び出し方。
688+
decreased = true;
689+
dealer.setLimit(1);
690+
resolveCanFinishFirstBatch();
691+
692+
await done;
693+
694+
expect(startedFirstBatch).toHaveLength(3);
695+
expect(maxConcurrentAfterDecrease).toBe(1);
696+
expect(dealer.limit).toBe(1);
697+
});
698+
699+
test('throws RangeError for non-positive-integer limits', () => {
700+
const dealer = new Dealer(createItems(1), { limit: 5 });
701+
expect(() => dealer.setLimit(0)).toThrow(RangeError);
702+
expect(() => dealer.setLimit(-1)).toThrow(RangeError);
703+
expect(() => dealer.setLimit(1.5)).toThrow(RangeError);
704+
expect(dealer.limit).toBe(5);
705+
});
706+
707+
test('after the dealer has finished, setLimit updates the stored limit without throwing or dispatching', async () => {
708+
const items = createItems(1);
709+
const dealer = new Dealer(items, { limit: 10 });
710+
711+
await dealer.setup(() => Promise.resolve(() => {}));
712+
await runDealer(dealer);
713+
714+
expect(() => dealer.setLimit(3)).not.toThrow();
715+
// #deal() 自体は #finished ガードで即 return する(新規ディスパッチは
716+
// 発生しない)が、#limit フィールドの更新はガードの影響を受けない
717+
expect(dealer.limit).toBe(3);
718+
});
719+
720+
test('a synchronous setLimit call from within a worker still respects the limit deterministically (re-entrant #deal() calls are ignored)', async () => {
721+
// worker 自身の同期区間(最初の await より前)から setLimit を呼ぶ
722+
// 稀なケースでも、#deal() の再入防止により外側のディスパッチループが
723+
// 一貫して最新の #limit を尊重する。increase 版は
724+
// 'onStart receives a controller...'(deal.spec.ts)で間接的に検証済み
725+
// なので、ここでは decrease 版のみ確認する。
726+
const items = createItems(3);
727+
const dealer = new Dealer(items, { limit: 3 });
728+
const processed: number[] = [];
729+
730+
await dealer.setup((_item, index) => {
731+
return Promise.resolve(async () => {
732+
processed.push(index);
733+
if (index === 0) {
734+
// この時点で #deal() はまだ while ループの最中(再入)
735+
dealer.setLimit(1);
736+
}
737+
await new Promise((r) => setTimeout(r, 1));
738+
});
739+
});
740+
741+
await runDealer(dealer);
742+
// 再入経路でも例外や取りこぼしなく全件完了する
743+
expect(processed).toHaveLength(3);
744+
expect(dealer.limit).toBe(1);
745+
});
746+
});
607747
});

0 commit comments

Comments
 (0)