@d-zero/dealer
v1.14.0
Published
A tool that provides an API and CLI for parallel processing of collections and sequential logging to standard output, plus a type-safe sequential pipeline builder that renders a live task-list TUI while running
Readme
@d-zero/dealer
コレクションを並列処理し、ログを順次出力する API(deal。並列数制御・キャンセル・進捗ヘッダー対応)と、
型安全な逐次パイプラインを構築するとタスクリストTUIが自動生成される API(TaskList)を提供する。
Installation
yarn add @d-zero/dealerUsage
import { deal } from '@d-zero/dealer';
await deal(
items,
(item, update, index, setLineHeader, push) => {
return async () => {
update(`item(${index}): processing`);
await item.start();
};
},
{ limit: 30 },
);setup コールバックは「初期化」を同期で行い、「実行関数」を返す形(並列度を超えた分はキューに入る)。push / unshift で実行中に新規アイテムを動的に追加可能。
キャンセル(AbortSignal)
const controller = new AbortController();
setTimeout(() => controller.abort(), 30_000);
await deal(items, setup, { limit: 10, signal: controller.signal });abort 時の挙動: 新規ワーカー起動を停止、実行中ワーカーは完了まで待機、push/unshift は無視。詳細は src/deal.ts / src/dealer.ts の JSDoc。
実行中の並列数変更・外部 Lanes の再利用
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)
import { TaskList } from '@d-zero/dealer';
const id = await TaskList.pipe('fetch', async () => fetchUser(userId))
.pipe('normalize', (user) => normalizeUser(user))
.pipe('save', async (user, ctx) => {
ctx.progress('writing to db...');
await db.save(user);
return user.id;
})
.run();.pipe() を連結するたびに前段の戻り値を受け取り型変換する新しいステップが追加され、run() を呼ぶと全ステップを最初から [ ] タスク名: 進捗メッセージ の形式で表示し、先頭から逐次実行しながら状態(pending/running/done/error)を更新する。あるステップが失敗すると即座に停止し、TaskListStepError で reject する(後続ステップは実行されない)。
- 各パイプラインインスタンスの
run()は1回のみ呼び出せる(再実行したい場合はTaskList.pipe()から新しく構築する) ctx.insertNext()で、実行中に同じ型のステップを直後へ動的に割り込み挿入できる- 詳細は
src/task-list-pipeline.ts/src/types.tsの JSDoc を参照
重要な制約
- **
interval遅延はアイテム開始の「直後・最初の出力前」**に実行される(順序に注意) unshiftは既存キューの先頭に割り込む(優先度の高い動的追加用、push との順序を理解する必要あり)Lanesを直接使う場合はusing宣言(Symbol.dispose)で自動解放する(leak 防止)。スコープと解放タイミングが一致しない場合のみclose()を直接呼ぶ(close()は deprecated)
これらの背景と実装は src/deal.ts / src/dealer.ts / src/lanes.ts の JSDoc を参照。
