-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathpipeline.ts
More file actions
293 lines (270 loc) · 11.5 KB
/
Copy pathpipeline.ts
File metadata and controls
293 lines (270 loc) · 11.5 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
import {
DEFAULT_SUMMARY_BATCH_SIZE,
DEFAULT_SUMMARY_INTERVAL_MS,
SUMMARY_RETRY_MAX_WAIT_MS,
} from "@/constants.js";
import { classifyQuotaError, parseRetryDelayMs } from "@/core/summarizer/quota.js";
import type { Article, RawArticle } from "@/types/article.js";
import type { ArticleSummary, Fetcher, Summarizer } from "@/types/ports.js";
import { describeError, logger } from "@/utils/logger.js";
/** Pipeline が必要とする実体依存。CLI 側で本物を、テストではモックを差し込む。 */
export interface PipelineDeps {
/** 取得元の Fetcher 一覧(Zenn / Qiita など)。 */
readonly fetchers: readonly Fetcher[];
/** 要約器(`--no-summary` のときは呼ばれない)。 */
readonly summarizer: Summarizer;
}
/** Pipeline の振る舞いを制御するオプション。 */
export interface PipelineOptions {
/** 要約を実行するか(`false` のとき summary には `RawArticle.excerpt` をそのまま使い、全件 `summarized:false`)。 */
readonly summarize: boolean;
/**
* 要約を試みる最大件数(トレンド上位から)。
*
* Gemini 無料枠の日次上限を踏まえた節約用。未指定なら全件試行し、上限到達(429)で自動停止する。
*/
readonly summaryLimit?: number;
/** 1 リクエストでまとめて要約する記事数。未指定は {@link DEFAULT_SUMMARY_BATCH_SIZE}。 */
readonly batchSize?: number;
/** 要約リクエスト間の待機ミリ秒(RPM 保護)。未指定は {@link DEFAULT_SUMMARY_INTERVAL_MS}。 */
readonly intervalMs?: number;
/** 待機関数(テスト注入用)。未指定は実 `setTimeout`。 */
readonly sleep?: (ms: number) => Promise<void>;
}
/**
* Pipeline 全体として処理を続けるべきでない失敗。
*
* 取得元の一部失敗は隔離して継続するが、全取得元が失敗した場合は空 JSON で
* 既存出力を上書きしないよう、CLI 側で systemError に変換する。
*/
export class PipelineError extends Error {
/**
* @param message - 利用者向けの説明
*/
constructor(message: string) {
super(message);
this.name = "PipelineError";
}
}
/**
* 取得 → 要約のオーケストレーション。
*
* - 取得は **ソース単位で並列・try/catch で失敗隔離**(1 サイト落ちても他は返る)
* - 件数のクランプは **Fetcher 自身の責務**(本文取得 API の往復を減らすため、
* Pipeline まで来た時点ではすでに上限内になっている前提)
* - 要約は **バッチ + スロットル**(Gemini 無料枠は 1 日あたりのリクエスト数が上限のため、
* 複数記事を 1 リクエストにまとめる)。日次上限(429/PerDay)に達したら以降のバッチを止め、
* 残りは `excerpt` を流用して `summarized:false` にする
* - 進捗・取得件数・要約件数は **stderr** へ
*
* @param deps - Fetcher 群と Summarizer
* @param options - 要約 ON/OFF・上限件数・スロットル間隔・待機関数
* @returns `Article` の配列(ソースの宣言順)。各要素は要約可否を `summarized` で持つ
*/
export async function runPipeline(
deps: PipelineDeps,
options: PipelineOptions,
): Promise<Article[]> {
const fetched = await fetchAll(deps.fetchers);
const collected = collect(fetched);
if (collected.allSourcesFailed) {
throw new PipelineError(
`すべての取得元で取得に失敗しました: ${collected.failedSources.join(", ")}`,
);
}
const rawArticles = collected.articles;
if (!options.summarize) {
return rawArticles.map((raw) => ({ ...raw, summary: raw.excerpt, summarized: false }));
}
logger.info(`↓ 要約中 (${rawArticles.length} 件)...`);
return summarizeArticles(rawArticles, deps.summarizer, options);
}
/**
* 全 Fetcher を並列で呼び、各々を「fetcher 自身とペアで」確定させる。
*
* `Promise.all` + 内側 try/catch を使うのは、rejection 時にも fetcher 情報を保持して
* 失敗ソース名をログに出せるようにするため。
*
* @param fetchers - 並列実行する Fetcher 群
* @returns 各 fetcher の成否ペア
*/
async function fetchAll(
fetchers: readonly Fetcher[],
): Promise<
ReadonlyArray<
| { readonly fetcher: Fetcher; readonly ok: true; readonly articles: readonly RawArticle[] }
| { readonly fetcher: Fetcher; readonly ok: false; readonly error: unknown }
>
> {
return Promise.all(
fetchers.map(async (fetcher) => {
try {
return { fetcher, ok: true as const, articles: await fetcher.fetch() };
} catch (error) {
return { fetcher, ok: false as const, error };
}
}),
);
}
/**
* 取得結果を集計して 1 つの配列に平坦化する。
*
* - 成功ソース: 取得件数を info で stderr に通知
* - 失敗ソース: warn を出して空件として扱う(停止しない)
*
* クランプは Fetcher の責務に集約しており、ここでは行わない(pipeline に来た時点で上限内)。
*
* @param fetched - {@link fetchAll} の結果
* @returns 全ソースを連結した RawArticle 配列と失敗元の集計
*/
function collect(
fetched: ReadonlyArray<
| { readonly fetcher: Fetcher; readonly ok: true; readonly articles: readonly RawArticle[] }
| { readonly fetcher: Fetcher; readonly ok: false; readonly error: unknown }
>,
): {
readonly articles: RawArticle[];
readonly failedSources: readonly string[];
readonly allSourcesFailed: boolean;
} {
const collected: RawArticle[] = [];
const summary: string[] = [];
const failedSources: string[] = [];
for (const entry of fetched) {
if (!entry.ok) {
logger.warn(`${entry.fetcher.source}: 取得に失敗: ${describeError(entry.error)}`);
summary.push(`${entry.fetcher.source} 0`);
failedSources.push(entry.fetcher.source);
continue;
}
collected.push(...entry.articles);
summary.push(`${entry.fetcher.source} ${entry.articles.length}`);
}
logger.info(`📰 取得: ${summary.join(" / ")} (合計 ${collected.length} 件)`);
return {
articles: collected,
failedSources,
allSourcesFailed: fetched.length > 0 && failedSources.length === fetched.length,
};
}
/** 1 バッチの要約結果。成功なら要約結果配列、失敗なら以降を停止すべきか(`stop`)を持つ。 */
type BatchOutcome =
| { readonly ok: true; readonly summaries: readonly ArticleSummary[] }
| { readonly ok: false; readonly stop: boolean };
/**
* 既定の待機関数(実時間スリープ)。
*
* @param ms - 待機ミリ秒
*/
function defaultSleep(ms: number): Promise<void> {
return new Promise((resolve) => setTimeout(resolve, ms));
}
/** 未要約の `Article`(`excerpt` を流用し `summarized:false`)を作る。 */
function toUnsummarized(raw: RawArticle): Article {
return { ...raw, summary: raw.excerpt, summarized: false };
}
/**
* 要約をバッチ単位で実行する。
*
* 先頭 `summaryLimit` 件を要約対象とし、`batchSize` ごとに 1 リクエストへまとめて要約する。
* RPM 保護のためリクエスト間に待機し、日次上限(429/PerDay)に達したら以降のバッチを止める。
* 上限超過・失敗したバッチや、`summaryLimit` を超えた記事は `excerpt` を流用して `summarized:false`。
*
* @param rawArticles - 取得済みの記事一覧(要約前)
* @param summarizer - 要約器(バッチ)
* @param options - 上限件数・バッチサイズ・スロットル間隔・待機関数を含む Pipeline オプション
* @returns `summarized` 付きの `Article` 配列(入力順を保つ)
*/
async function summarizeArticles(
rawArticles: readonly RawArticle[],
summarizer: Summarizer,
options: PipelineOptions,
): Promise<Article[]> {
const intervalMs = options.intervalMs ?? DEFAULT_SUMMARY_INTERVAL_MS;
const sleep = options.sleep ?? defaultSleep;
const batchSize = Math.max(1, options.batchSize ?? DEFAULT_SUMMARY_BATCH_SIZE);
const attemptCount =
options.summaryLimit === undefined
? rawArticles.length
: Math.min(options.summaryLimit, rawArticles.length);
const articles: Article[] = [];
let quotaExhausted = false;
let calledOnce = false;
for (let start = 0; start < attemptCount; start += batchSize) {
const batch = rawArticles.slice(start, Math.min(start + batchSize, attemptCount));
if (quotaExhausted) {
for (const raw of batch) articles.push(toUnsummarized(raw));
continue;
}
// 実リクエストの 2 回目以降は間隔を空ける(RPM 保護)。
if (calledOnce) await sleep(intervalMs);
calledOnce = true;
const outcome = await summarizeBatchOnce(summarizer, batch, sleep);
if (outcome.ok) {
batch.forEach((raw, j) => {
const result = outcome.summaries[j];
articles.push({
...raw,
summary: result?.summary ?? raw.excerpt,
...(result?.titleTranslated ? { titleTranslated: result.titleTranslated } : {}),
summarized: true,
});
});
} else {
if (outcome.stop) {
quotaExhausted = true;
logger.warn("無料枠の上限に達したため、以降は要約せず本文を表示します");
}
for (const raw of batch) articles.push(toUnsummarized(raw));
}
}
// summaryLimit を超えた残りは未要約。
for (const raw of rawArticles.slice(attemptCount)) articles.push(toUnsummarized(raw));
const summarizedCount = articles.reduce((count, a) => count + (a.summarized ? 1 : 0), 0);
logger.info(`要約: 成功 ${summarizedCount} / 未要約 ${articles.length - summarizedCount} 件`);
return articles;
}
/**
* 1 バッチを要約する。分次レート超過(429/PerMinute)なら retryDelay だけ待って 1 回だけ再試行する。
*
* @param summarizer - 要約器(バッチ)
* @param batch - このバッチの記事
* @param sleep - 待機関数(再試行の待ちに使う)
* @returns 成功なら要約文配列、失敗なら以降を止めるべきか(`stop`)を含む結果
*/
async function summarizeBatchOnce(
summarizer: Summarizer,
batch: readonly RawArticle[],
sleep: (ms: number) => Promise<void>,
): Promise<BatchOutcome> {
try {
const summaries = await summarizer.summarizeBatch(batch);
if (summaries.length !== batch.length) {
logger.warn(`要約バッチの件数不一致(${summaries.length}/${batch.length})、本文を表示`);
return { ok: false, stop: false };
}
return { ok: true, summaries };
} catch (error) {
const kind = classifyQuotaError(error);
if (kind === undefined) {
// クォータ以外の失敗(解析失敗・一時障害など)→ このバッチだけ未要約にして継続。
logger.warn(`要約バッチ失敗、本文を表示: ${describeError(error)}`);
return { ok: false, stop: false };
}
if (kind === "rate") {
const waitMs = parseRetryDelayMs(error, SUMMARY_RETRY_MAX_WAIT_MS);
logger.warn(`レート制限。${(waitMs / 1000).toFixed(0)}s 待って再試行(${batch.length} 件)`);
await sleep(waitMs);
try {
const summaries = await summarizer.summarizeBatch(batch);
if (summaries.length === batch.length) return { ok: true, summaries };
return { ok: false, stop: false };
} catch (retryError) {
logger.warn(`再試行も失敗: ${describeError(retryError)}`);
return { ok: false, stop: true };
}
}
// kind === "daily": 日次上限 → 以降を止める。
return { ok: false, stop: true };
}
}