From e454848aa07b60712a0797a4ae11ce7b386359a2 Mon Sep 17 00:00:00 2001 From: Yusuke Hirao Date: Tue, 8 Sep 2026 17:41:27 +0900 Subject: [PATCH 1/2] feat(dealer): support runtime concurrency changes and injected Lanes Add Dealer#setLimit()/#limit to change worker concurrency mid-run, with a reentrancy guard so a synchronous call from within a worker's own start() still dispatches deterministically. Add Lanes#footer() (and clear({footer})) as a fixed line below the log lanes, and DealOptions .lanes/.onStart so a caller can reuse an existing Lanes instance and get a DealController to call setLimit() from an external event (e.g. a CLI reading stdin while a deal() is in flight). --- packages/@d-zero/dealer/README.md | 19 +++ packages/@d-zero/dealer/src/deal.spec.ts | 71 +++++++++++ packages/@d-zero/dealer/src/deal.ts | 26 +++- packages/@d-zero/dealer/src/dealer.spec.ts | 140 +++++++++++++++++++++ packages/@d-zero/dealer/src/dealer.ts | 118 ++++++++++++----- packages/@d-zero/dealer/src/index.ts | 1 + packages/@d-zero/dealer/src/lanes.spec.ts | 61 +++++++++ packages/@d-zero/dealer/src/lanes.ts | 44 ++++++- packages/@d-zero/dealer/src/types.ts | 31 +++++ 9 files changed, 473 insertions(+), 38 deletions(-) diff --git a/packages/@d-zero/dealer/README.md b/packages/@d-zero/dealer/README.md index 4e44e626..5febfbc3 100644 --- a/packages/@d-zero/dealer/README.md +++ b/packages/@d-zero/dealer/README.md @@ -39,6 +39,25 @@ 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; + +await deal(items, setup, { + limit: 10, + lanes, // 呼び出し元が生成した Lanes を使い回す(deal() は生成も破棄もしない) + onStart: (c) => { + controller = c; + }, +}); +``` + +`lanes` を渡すと `deal()` は自前で `Lanes` を作らず、渡されたインスタンスの生成・破棄は呼び出し元の責任になる。`onStart` は `dealer.play()` 直前に一度呼ばれ、`controller.setLimit(n)` で実行中に並列数を変更できる(`Lanes#footer(text)` と組み合わせれば、外部からの入力受付 UI を並列レーンの下に固定表示できる)。 + ## Sequential Pipeline(`TaskList`) ```ts diff --git a/packages/@d-zero/dealer/src/deal.spec.ts b/packages/@d-zero/dealer/src/deal.spec.ts index c997dcb9..5aafbead 100644 --- a/packages/@d-zero/dealer/src/deal.spec.ts +++ b/packages/@d-zero/dealer/src/deal.spec.ts @@ -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'; /** * @@ -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(); + const { promise: canFinishFirst, resolve: resolveCanFinishFirst } = + Promise.withResolvers(); + 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); + }); }); diff --git a/packages/@d-zero/dealer/src/deal.ts b/packages/@d-zero/dealer/src/deal.ts index 6ee54287..aef58132 100644 --- a/packages/@d-zero/dealer/src/deal.ts +++ b/packages/@d-zero/dealer/src/deal.ts @@ -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'; @@ -18,6 +19,19 @@ export type DealOptions = DealerOptions & 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; }; /** @@ -98,7 +112,11 @@ export async function deal( 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) => { @@ -137,6 +155,12 @@ export async function deal( // 実行される必要がある)ため、ここは `await` で完了を待ってからスコープを抜ける。 const { promise, resolve } = Promise.withResolvers(); dealer.finish(resolve); + options?.onStart?.({ + get limit() { + return dealer.limit; + }, + setLimit: (limit) => dealer.setLimit(limit), + }); dealer.play(); await promise; } diff --git a/packages/@d-zero/dealer/src/dealer.spec.ts b/packages/@d-zero/dealer/src/dealer.spec.ts index bd0a73b3..2eb42d15 100644 --- a/packages/@d-zero/dealer/src/dealer.spec.ts +++ b/packages/@d-zero/dealer/src/dealer.spec.ts @@ -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(); + const { promise: canFinishFirst, resolve: resolveCanFinishFirst } = + Promise.withResolvers(); + 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(); + const { promise: canFinishFirstBatch, resolve: resolveCanFinishFirstBatch } = + Promise.withResolvers(); + 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); + }); + }); }); diff --git a/packages/@d-zero/dealer/src/dealer.ts b/packages/@d-zero/dealer/src/dealer.ts index a5ea7c4b..f650fa8c 100644 --- a/packages/@d-zero/dealer/src/dealer.ts +++ b/packages/@d-zero/dealer/src/dealer.ts @@ -33,6 +33,7 @@ export interface DealerOptions { * @template T - 処理対象アイテムの型(WeakKey 制約) */ export class Dealer { + #dealing = false; #debug: (log: string) => void = () => {}; #done = new WeakSet(); #doneCount = 0; @@ -50,6 +51,12 @@ export class Dealer { #starts = new WeakMap Promise>(); #workers = new Set(); + /** + * 現在の同時実行ワーカー数の上限。 + */ + get limit(): number { + return this.#limit; + } constructor(items: readonly T[], options?: DealerOptions) { this.#items = [...items]; this.#limit = options?.limit ?? 10; @@ -102,6 +109,30 @@ export class Dealer { await this.#enqueue(items, false); } + /** + * 同時実行ワーカー数の上限を実行中に変更する。 + * 増加時は空きスロットを即座に充填する。減少時は新規ワーカーの起動を + * 絞るだけで、既に実行中のワーカーを中断しない — + * `#workers.size` が新しい上限を下回るまで自然に収束する。 + * + * 呼び出し元が別イベント(タイマー・外部入力等)から呼ぶ通常のケースでは + * 上記の効果が即座に反映される。ワーカー自身の `start()` の同期区間 + * (最初の `await` より前)から呼んだ場合も安全(`#deal()` の再入防止に + * より、その時点の `#limit` を一貫して尊重する)だが、増加時に空きスロット + * を埋めるのはディスパッチ元の `#deal()` 呼び出し自身の残りイテレーションに + * なる — 現在ディスパッチ中のワーカー自身がまだ起動されていない同バッチの + * 他アイテムより先に完了した場合、それらのアイテムは新しい上限を下回るまで + * 起動が遅れることがある。 + * @param limit - 新しい上限(1以上の整数) + * @throws {RangeError} `limit` が1以上の整数でない場合 + */ + setLimit(limit: number) { + if (!Number.isInteger(limit) || limit < 1) { + throw new RangeError(`limit must be an integer >= 1, got ${limit}`); + } + this.#limit = limit; + this.#deal(); + } /** * 各アイテムの初期化関数を設定する。 * {@link play} を呼ぶ前に必ず呼び出すこと。 @@ -135,51 +166,72 @@ export class Dealer { await this.#enqueue(items, true); } + /** + * ディスパッチ・完了判定を行う中核ループ。 + * + * 再入防止(`#dealing`): `start()` の同期区間(最初の `await` まで)から + * {@link setLimit} が呼ばれると、`#deal()` が自分自身の `while` ループの + * 内側から再帰的に呼ばれる。ガードなしだと、ネストした呼び出しが増えた + * 枠を先取りしたり、外側のループが変更後の `#limit` で次のイテレーション + * を評価したりする順序がタイミング依存になり、どちらが何件ディスパッチ + * するかが不定になる。`#dealing` が立っている間の再入は即座に無視し、 + * 外側の `while` ループ自身が次のイテレーションで最新の `#limit`/ + * `#workers.size` を読み直して続行する — 増加時は外側ループが残り枠を + * 埋め、減少時は外側ループがそこで打ち切る。結果はどちらも「その時点の + * `#limit` を尊重する」という同じ規則の一貫した適用になり、呼び出しが + * ワーカー自身の同期区間から来たか外部からの非同期呼び出しかによらず + * 決定的になる。 + */ #deal() { - if (this.#finished) { + if (this.#finished || this.#dealing) { return; } - const total = this.#items.length; - this.#debug(`Done: ${this.#doneCount}/${total} (Limit: ${this.#limit})`); - this.#progress( - total === 0 ? 0 : this.#doneCount / total, - this.#doneCount, - total, - this.#limit, - ); + this.#dealing = true; + try { + const total = this.#items.length; + this.#debug(`Done: ${this.#doneCount}/${total} (Limit: ${this.#limit})`); + this.#progress( + total === 0 ? 0 : this.#doneCount / total, + this.#doneCount, + total, + this.#limit, + ); - if (this.#doneCount === total && this.#pendingInitCount === 0) { - this.#finished = true; - this.#finish(); - return; - } - - if (this.#signal?.aborted) { - if (this.#workers.size === 0) { + if (this.#doneCount === total && this.#pendingInitCount === 0) { this.#finished = true; this.#finish(); + return; } - return; - } - while (this.#workers.size < this.#limit) { - const worker = this.#draw(); - if (!worker) { + if (this.#signal?.aborted) { + if (this.#workers.size === 0) { + this.#finished = true; + this.#finish(); + } return; } - this.#workers.add(worker); - const start = this.#starts.get(worker); - if (!start) { - throw new Error(`Didn't have a starting function`); - } + while (this.#workers.size < this.#limit) { + const worker = this.#draw(); + if (!worker) { + return; + } + + this.#workers.add(worker); + const start = this.#starts.get(worker); + if (!start) { + throw new Error(`Didn't have a starting function`); + } - void start().then(() => { - this.#workers.delete(worker); - this.#done.add(worker); - this.#doneCount++; - this.#deal(); - }); + void start().then(() => { + this.#workers.delete(worker); + this.#done.add(worker); + this.#doneCount++; + this.#deal(); + }); + } + } finally { + this.#dealing = false; } } #draw() { diff --git a/packages/@d-zero/dealer/src/index.ts b/packages/@d-zero/dealer/src/index.ts index 579c815d..b33bcd0a 100644 --- a/packages/@d-zero/dealer/src/index.ts +++ b/packages/@d-zero/dealer/src/index.ts @@ -3,6 +3,7 @@ export { deal } from './deal.js'; export { Dealer } from './dealer.js'; export { Lanes } from './lanes.js'; export type { + DealController, StepContext, StepFn, TaskListRunOptions, diff --git a/packages/@d-zero/dealer/src/lanes.spec.ts b/packages/@d-zero/dealer/src/lanes.spec.ts index ff36fb97..e311371c 100644 --- a/packages/@d-zero/dealer/src/lanes.spec.ts +++ b/packages/@d-zero/dealer/src/lanes.spec.ts @@ -100,6 +100,67 @@ describe('Lanes verbose update', () => { }); }); +describe('Lanes footer', () => { + test('footer() renders after the lane logs and after the header', () => { + const collector = makeStreamCollector(); + using lanes = new Lanes({ stream: collector.stream }); + + lanes.header('My Header'); + lanes.update(0, 'lane log'); + lanes.footer('> input line'); + // header()/update()/footer() 呼び出しはタイマーが既にペンディング中なら + // #stack を更新するだけで同期描画しない。resize は #write() を直接叩き、 + // 現在の #stack を強制的に同期描画する(countdown lifetime テストと同じ手法)。 + collector.stream.emit('resize'); + + const painted = collector.read(); + const headerIndex = painted.indexOf('My Header'); + const logIndex = painted.lastIndexOf('lane log'); + const footerIndex = painted.lastIndexOf('> input line'); + + expect(headerIndex).toBeGreaterThanOrEqual(0); + expect(logIndex).toBeGreaterThan(headerIndex); + expect(footerIndex).toBeGreaterThan(logIndex); + }); + + test('footer() supports multi-line text', () => { + const collector = makeStreamCollector(); + using lanes = new Lanes({ stream: collector.stream }); + + lanes.footer('status line\n> input'); + + const painted = collector.read(); + expect(painted).toContain('status line'); + expect(painted).toContain('> input'); + }); + + test('clear({ footer: true }) removes the footer, clear() alone keeps it', () => { + const collector = makeStreamCollector(); + using lanes = new Lanes({ stream: collector.stream }); + + lanes.footer('> input line'); + + lanes.clear(); + const markAfterClear = collector.read().length; + collector.stream.emit('resize'); + expect(collector.read().slice(markAfterClear)).toContain('> input line'); + + lanes.clear({ footer: true }); + const markAfterClearFooter = collector.read().length; + collector.stream.emit('resize'); + expect(collector.read().slice(markAfterClearFooter)).not.toContain('> input line'); + }); + + test('footer() is a no-op in verbose mode', () => { + using lanes = new Lanes({ verbose: true }); + stdoutWriteSpy.mockClear(); + + lanes.footer('> input line'); + + expect(stdoutWriteSpy).not.toHaveBeenCalled(); + }); +}); + describe('Lanes countdown lifetime', () => { test('a lane-scoped countdown id restarts from its full duration on the next item', () => { vi.useFakeTimers(); diff --git a/packages/@d-zero/dealer/src/lanes.ts b/packages/@d-zero/dealer/src/lanes.ts index 1e55a9d4..84643ccd 100644 --- a/packages/@d-zero/dealer/src/lanes.ts +++ b/packages/@d-zero/dealer/src/lanes.ts @@ -7,6 +7,16 @@ const RESET = '\u001B[0m'; type Log = readonly [id: number, message: string]; type SortFunc = (a: Log, b: Log) => number; +/** + * header / footer 共通の整形処理。複数行テキストを改行で分割し、 + * 各行を `RESET` で囲む(ログ行に紛れて前の行の色が漏れ継続しないよう)。 + * @param text - `header()` / `footer()` に渡された生テキスト + * @returns `Display.write()` にそのまま渡せる行の配列 + */ +function formatFixedLines(text: string): string[] { + return text.split('\n').map((line) => `${RESET}${line}${RESET}`); +} + /** * {@link Lanes} のコンストラクタオプション。 */ @@ -30,6 +40,7 @@ export type LanesOptions = { */ export class Lanes { #display: Display; + #footerText?: string; #header?: string; #indent = ''; #logs = new Map(); @@ -65,8 +76,9 @@ export class Lanes { * すべてのログをクリアする。verbose モードでは何もしない。 * @param options - クリアオプション * @param options.header + * @param options.footer */ - clear(options?: { header?: boolean }) { + clear(options?: { header?: boolean; footer?: boolean }) { if (this.#verbose) { return; } @@ -77,6 +89,10 @@ export class Lanes { this.#header = undefined; } + if (options?.footer) { + this.#footerText = undefined; + } + this.write(); } /** @@ -100,6 +116,25 @@ export class Lanes { this.#logs.delete(id); this.write(); } + /** + * フッターテキストを設定する。{@link header} と対称で、全レーンの + * ログの下に固定表示される。常時表示の入力行など、ログの再描画に + * 巻き込まれずに末尾へ固定したい行に使う。 + * + * verbose モードでは何もしない — verbose には上書きされる単一フレームが + * 存在せず、{@link update} 呼び出しのたびにフッターが追記出力され続けると + * スパムになるため。 + * @param text - フッターとして表示する文字列 + */ + footer(text: string) { + if (this.#verbose) { + return; + } + + this.#footerText = text; + this.write(); + } + /** * ヘッダーテキストを設定する。 * @param text - ヘッダーとして表示する文字列 @@ -154,9 +189,10 @@ export class Lanes { logs.sort(this.#sort); const messages = logs.map(([, message]) => `${this.#indent}${message}`); if (this.#header) { - messages.unshift( - ...this.#header.split('\n').map((line) => `${RESET}${line}${RESET}`), - ); + messages.unshift(...formatFixedLines(this.#header)); + } + if (this.#footerText) { + messages.push(...formatFixedLines(this.#footerText)); } this.#display.write(...messages); } diff --git a/packages/@d-zero/dealer/src/types.ts b/packages/@d-zero/dealer/src/types.ts index 8d70ed76..ec422ad2 100644 --- a/packages/@d-zero/dealer/src/types.ts +++ b/packages/@d-zero/dealer/src/types.ts @@ -21,6 +21,37 @@ export interface ProcessInitializer { (process: T, index: number): Promise<() => Promise | void>; } +/** + * {@link deal} 実行中に同時実行数を操作するためのハンドル。 + * `options.onStart` で `dealer.play()` 直前に一度だけ渡される。 + * + * `deal()` が返す Promise は全アイテムの処理完了まで解決しないため、 + * `controller` は `await deal(...)` の**外側**、`deal()` 実行中に発火する + * 別イベント(CLI のキー入力、タイマー等)から使う。`await deal(...)` の + * 後で呼んでも例外にはならないが、その時点で全ワーカーは完了済みのため + * 並列数の変更対象がなく無意味になる。 + * @example + * ```ts + * let controller: DealController | undefined; + * const run = deal(items, setup, { + * onStart: (c) => { controller = c; }, + * }); + * // deal() 実行中(run が解決する前)に、別イベントから並列数を変更する + * onExternalCommand((newLimit) => controller?.setLimit(newLimit)); + * await run; + * ``` + */ +export interface DealController { + /** 現在の同時実行ワーカー数の上限。 */ + readonly limit: number; + /** + * 同時実行ワーカー数の上限を変更する。 + * @param limit - 新しい上限(1以上の整数) + * @throws {RangeError} `limit` が1以上の整数でない場合 + */ + setLimit(limit: number): void; +} + /** * {@link TaskListPipeline} の各ステップの実行状態。 * - `pending`: 未実行 From f6496139bc8607b92c2e8b531f76995eb8fbb972 Mon Sep 17 00:00:00 2001 From: Yusuke Hirao Date: Tue, 8 Sep 2026 17:50:49 +0900 Subject: [PATCH 2/2] docs(dealer): fix README example showing setLimit() called before await deal() resolves --- packages/@d-zero/dealer/README.md | 9 +++++++-- 1 file changed, 7 insertions(+), 2 deletions(-) diff --git a/packages/@d-zero/dealer/README.md b/packages/@d-zero/dealer/README.md index 5febfbc3..fa8bc689 100644 --- a/packages/@d-zero/dealer/README.md +++ b/packages/@d-zero/dealer/README.md @@ -47,16 +47,21 @@ import { deal, Lanes } from '@d-zero/dealer'; const lanes = new Lanes({ stream: process.stderr }); let controller: DealController | undefined; -await deal(items, setup, { +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)` で実行中に並列数を変更できる(`Lanes#footer(text)` と組み合わせれば、外部からの入力受付 UI を並列レーンの下に固定表示できる)。 +`lanes` を渡すと `deal()` は自前で `Lanes` を作らず、渡されたインスタンスの生成・破棄は呼び出し元の責任になる。`onStart` は `dealer.play()` 直前に一度呼ばれ、`controller.setLimit(n)` を **`await deal(...)` が解決する前に** 呼ぶことで実行中に並列数を変更できる(`Lanes#footer(text)` と組み合わせれば、外部からの入力受付 UI を並列レーンの下に固定表示できる)。 ## Sequential Pipeline(`TaskList`)