Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
24 changes: 24 additions & 0 deletions packages/@d-zero/dealer/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -39,6 +39,30 @@ await deal(items, setup, { limit: 10, signal: controller.signal });

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

### 実行中の並列数変更・外部 `Lanes` の再利用

```ts
import { deal, Lanes } from '@d-zero/dealer';

const lanes = new Lanes({ stream: process.stderr });
let controller: DealController | undefined;

const run = deal(items, setup, {
limit: 10,
lanes, // 呼び出し元が生成した Lanes を使い回す(deal() は生成も破棄もしない)
onStart: (c) => {
controller = c;
},
});

// deal() 実行中(run が解決する前)に、別イベントから並列数を変更する
onExternalCommand((newLimit) => controller?.setLimit(newLimit));

await run;
```

`lanes` を渡すと `deal()` は自前で `Lanes` を作らず、渡されたインスタンスの生成・破棄は呼び出し元の責任になる。`onStart` は `dealer.play()` 直前に一度呼ばれ、`controller.setLimit(n)` を **`await deal(...)` が解決する前に** 呼ぶことで実行中に並列数を変更できる(`Lanes#footer(text)` と組み合わせれば、外部からの入力受付 UI を並列レーンの下に固定表示できる)。

## Sequential Pipeline(`TaskList`)

```ts
Expand Down
71 changes: 71 additions & 0 deletions packages/@d-zero/dealer/src/deal.spec.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,9 @@
import type { DealController } from './types.js';

import { describe, test, expect, vi } from 'vitest';

import { deal } from './deal.js';
import { Lanes } from './lanes.js';

/**
*
Expand Down Expand Up @@ -152,4 +155,72 @@ describe('deal', () => {

stdoutWriteSpy.mockRestore();
});

test('reuses an injected Lanes instead of creating its own, and does not dispose it', async () => {
const resizeBefore = process.stdout.listenerCount('resize');
const lanes = new Lanes({ verbose: true });
// injected の場合、deal() が独自の Lanes を作らないので resize リスナーは増えない
expect(process.stdout.listenerCount('resize')).toBe(resizeBefore + 1);

await deal(
createItems(2),
(_process, update, index) => {
return () => {
update(`item ${index}`);
};
},
{ limit: 10, lanes },
);

// deal() 完了後も呼び出し元の Lanes は破棄されず生きている
expect(process.stdout.listenerCount('resize')).toBe(resizeBefore + 1);
lanes[Symbol.dispose]();
expect(process.stdout.listenerCount('resize')).toBe(resizeBefore);
});

test('onStart receives a controller before play(), and setLimit affects the header limit', async () => {
const limits: number[] = [];
let controller: DealController | undefined;
const { promise: firstStarted, resolve: resolveFirstStarted } =
Promise.withResolvers<void>();
const { promise: canFinishFirst, resolve: resolveCanFinishFirst } =
Promise.withResolvers<void>();
let firstCall = true;

const run = deal(
createItems(3),
() => {
return async () => {
if (firstCall) {
firstCall = false;
resolveFirstStarted();
await canFinishFirst;
}
};
},
{
limit: 1,
verbose: true,
onStart: (c) => {
controller = c;
},
header: (_progress, _done, _total, limit) => {
limits.push(limit);
return `limit: ${limit}`;
},
},
);

await firstStarted;
expect(controller).toBeDefined();
expect(controller?.limit).toBe(1);

controller?.setLimit(3);
expect(controller?.limit).toBe(3);
resolveCanFinishFirst();

await run;

expect(limits).toContain(3);
});
});
26 changes: 25 additions & 1 deletion packages/@d-zero/dealer/src/deal.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
import type { DealerOptions } from './dealer.js';
import type { LanesOptions } from './lanes.js';
import type { DealController } from './types.js';
import type { DelayOptions } from '@d-zero/shared/delay';

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

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

if (options?.header) {
dealer.progress((progress, done, total, limit) => {
Expand Down Expand Up @@ -137,6 +155,12 @@ export async function deal<T extends WeakKey>(
// 実行される必要がある)ため、ここは `await` で完了を待ってからスコープを抜ける。
const { promise, resolve } = Promise.withResolvers<void>();
dealer.finish(resolve);
options?.onStart?.({
get limit() {
return dealer.limit;
},
setLimit: (limit) => dealer.setLimit(limit),
});
dealer.play();
await promise;
}
140 changes: 140 additions & 0 deletions packages/@d-zero/dealer/src/dealer.spec.ts
Original file line number Diff line number Diff line change
Expand Up @@ -604,4 +604,144 @@ describe('Dealer', () => {

await runDealer(dealer);
});

describe('setLimit', () => {
test('increasing the limit immediately fills newly available slots', async () => {
const items = createItems(4);
let maxConcurrent = 0;
let currentConcurrent = 0;
const dealer = new Dealer(items, { limit: 1 });
const { promise: firstStarted, resolve: resolveFirstStarted } =
Promise.withResolvers<void>();
const { promise: canFinishFirst, resolve: resolveCanFinishFirst } =
Promise.withResolvers<void>();
let firstCall = true;

await dealer.setup(() => {
return Promise.resolve(async () => {
currentConcurrent++;
maxConcurrent = Math.max(maxConcurrent, currentConcurrent);
if (firstCall) {
firstCall = false;
resolveFirstStarted();
await canFinishFirst;
}
currentConcurrent--;
});
});

const done = runDealer(dealer);
await firstStarted;
// limit 1 のあいだは1件しか動いていないはず
expect(maxConcurrent).toBe(1);

dealer.setLimit(4);
resolveCanFinishFirst();
await done;

expect(maxConcurrent).toBeGreaterThan(1);
expect(dealer.limit).toBe(4);
});

test('decreasing the limit lets already-started workers finish, but throttles concurrency for items dispatched afterward', async () => {
// limit 3・5件: 最初の3件(index 0,1,2)が同時ディスパッチされる。
// それらが実行中のうちに limit を 1 へ落とし、(a) 実行中の3件は
// 中断されず全件完了すること、(b) 減少後に新規ディスパッチされる
// 残り2件(index 3,4)は同時に1件までしか動かないことを検証する。
const items = createItems(5);
const dealer = new Dealer(items, { limit: 3 });
const startedFirstBatch: number[] = [];
const { promise: firstBatchStarted, resolve: resolveFirstBatchStarted } =
Promise.withResolvers<void>();
const { promise: canFinishFirstBatch, resolve: resolveCanFinishFirstBatch } =
Promise.withResolvers<void>();
let decreased = false;
let concurrentAfterDecrease = 0;
let maxConcurrentAfterDecrease = 0;

await dealer.setup((_item, index) => {
return Promise.resolve(async () => {
if (!decreased) {
startedFirstBatch.push(index);
if (startedFirstBatch.length === 3) {
resolveFirstBatchStarted();
}
await canFinishFirstBatch;
return;
}
concurrentAfterDecrease++;
maxConcurrentAfterDecrease = Math.max(
maxConcurrentAfterDecrease,
concurrentAfterDecrease,
);
await new Promise((r) => setTimeout(r, 5));
concurrentAfterDecrease--;
});
});

const done = runDealer(dealer);
await firstBatchStarted;
expect(startedFirstBatch).toHaveLength(3);

// setLimit はテストの非同期フロー(=ワーカー自身の同期区間の外)から
// 呼ぶ。これは実運用(外部入力による並列数変更)と同じ呼び出し方。
decreased = true;
dealer.setLimit(1);
resolveCanFinishFirstBatch();

await done;

expect(startedFirstBatch).toHaveLength(3);
expect(maxConcurrentAfterDecrease).toBe(1);
expect(dealer.limit).toBe(1);
});

test('throws RangeError for non-positive-integer limits', () => {
const dealer = new Dealer(createItems(1), { limit: 5 });
expect(() => dealer.setLimit(0)).toThrow(RangeError);
expect(() => dealer.setLimit(-1)).toThrow(RangeError);
expect(() => dealer.setLimit(1.5)).toThrow(RangeError);
expect(dealer.limit).toBe(5);
});

test('after the dealer has finished, setLimit updates the stored limit without throwing or dispatching', async () => {
const items = createItems(1);
const dealer = new Dealer(items, { limit: 10 });

await dealer.setup(() => Promise.resolve(() => {}));
await runDealer(dealer);

expect(() => dealer.setLimit(3)).not.toThrow();
// #deal() 自体は #finished ガードで即 return する(新規ディスパッチは
// 発生しない)が、#limit フィールドの更新はガードの影響を受けない
expect(dealer.limit).toBe(3);
});

test('a synchronous setLimit call from within a worker still respects the limit deterministically (re-entrant #deal() calls are ignored)', async () => {
// worker 自身の同期区間(最初の await より前)から setLimit を呼ぶ
// 稀なケースでも、#deal() の再入防止により外側のディスパッチループが
// 一貫して最新の #limit を尊重する。increase 版は
// 'onStart receives a controller...'(deal.spec.ts)で間接的に検証済み
// なので、ここでは decrease 版のみ確認する。
const items = createItems(3);
const dealer = new Dealer(items, { limit: 3 });
const processed: number[] = [];

await dealer.setup((_item, index) => {
return Promise.resolve(async () => {
processed.push(index);
if (index === 0) {
// この時点で #deal() はまだ while ループの最中(再入)
dealer.setLimit(1);
}
await new Promise((r) => setTimeout(r, 1));
});
});

await runDealer(dealer);
// 再入経路でも例外や取りこぼしなく全件完了する
expect(processed).toHaveLength(3);
expect(dealer.limit).toBe(1);
});
});
});
Loading
Loading