返回 CodeWhale
community-agent-tasks.ts
根目录 / web / lib / community-agent-tasks.ts
1 import type { DraftClaimLockNamespace, DraftClaimLockStub } from "./draft-claim-lock";
2 import { OUTBOUND_TIMEOUT_MS } from "@/lib/bounded-body";
3 import { fetchFeed, fetchRepoStats } from "@/lib/github";
4 import { curate } from "@/lib/deepseek";
5 import { putDispatchWithKv } from "@/lib/kv";
6 import {
7 agentChat,
8 TRIAGE_PROMPT,
9 PR_REVIEW_PROMPT,
10 STALE_PROMPT,
11 DUPES_PROMPT,
12 DIGEST_PROMPT,
13 saveDraft,
14 hasFreshDraft,
15 getDraftResolution,
16 digestRecordKey,
17 renderDigestBody,
18 logUsage,
19 type WeeklyDigestRecord,
20 type AgentDraft,
21 type DeepSeekEnv,
22 } from "@/lib/community-agent";
23
24 export interface AgentEnv {
25 DRAFT_CLAIM_LOCK?: DraftClaimLockNamespace;
26 CURATED_KV?: {
27 get(k: string): Promise<string | null>;
28 put(k: string, v: string, o?: { expirationTtl?: number }): Promise<void>;
29 list(o?: { prefix?: string; limit?: number }): Promise<{ keys: { name: string }[] }>;
30 delete(key: string): Promise<void>;
31 };
32 DEEPSEEK_API_KEY?: string;
33 DEEPSEEK_BASE_URL?: string;
34 DEEPSEEK_MODEL?: string;
35 GITHUB_TOKEN?: string;
36 CRON_SECRET?: string;
37 GITHUB_REPO?: string;
38 MAINTAINER_TOKEN?: string;
39 MAINTAINER_GITHUB_PAT?: string;
40 }
41
42 const CRON_STATUS_TTL = 60 * 60 * 24 * 14;
43
44 function dsEnv(env: AgentEnv): DeepSeekEnv {
45 return {
46 baseUrl: env.DEEPSEEK_BASE_URL ?? process.env.DEEPSEEK_BASE_URL,
47 model: env.DEEPSEEK_MODEL ?? process.env.DEEPSEEK_MODEL,
48 };
49 }
50
51 async function generateCurate(env: AgentEnv): Promise<Record<string, unknown>> {
52 if (!env.DEEPSEEK_API_KEY) {
53 return { skipped: true, reason: "DEEPSEEK_API_KEY not set" };
54 }
55 try {
56 const [stats, feed] = await Promise.all([
57 fetchRepoStats(env.GITHUB_TOKEN),
58 fetchFeed(env.GITHUB_TOKEN, 30),
59 ]);
60 const dispatch = await curate(env.DEEPSEEK_API_KEY, stats, feed, dsEnv(env));
61 await putDispatchWithKv(env.CURATED_KV, dispatch);
62 await env.CURATED_KV?.put(
63 "cron:curate:last",
64 JSON.stringify({
65 ok: true,
66 generatedAt: dispatch.generatedAt,
67 headline: dispatch.headline,
68 }),
69 { expirationTtl: CRON_STATUS_TTL }
70 );
71 return { ok: true, headline: dispatch.headline, stored: env.CURATED_KV ? "kv" : "memory" };
72 } catch (e) {
73 const error = String(e);
74 await env.CURATED_KV?.put(
75 "cron:curate:last",
76 JSON.stringify({
77 ok: false,
78 generatedAt: new Date().toISOString(),
79 error,
80 }),
81 { expirationTtl: CRON_STATUS_TTL }
82 );
83 return { ok: false, error };
84 }
85 }
86
87 async function generateTriage(env: AgentEnv): Promise<Record<string, unknown>> {
88 const repo = env.GITHUB_REPO ?? "codewhale-hq/CodeWhale";
89 try {
90 const res = await fetch(
91 `https://api.github.com/repos/${repo}/issues?state=open&sort=created&direction=desc&per_page=30`,
92 {
93 signal: AbortSignal.timeout(OUTBOUND_TIMEOUT_MS),
94 headers: {
95 Accept: "application/vnd.github+json",
96 "X-GitHub-Api-Version": "2022-11-28",
97 "User-Agent": "codewhale-web",
98 ...(env.GITHUB_TOKEN ? { Authorization: `Bearer ${env.GITHUB_TOKEN}` } : {}),
99 },
100 }
101 );
102 const issues = (await res.json()) as { number: number; title: string; body?: string; updated_at: string; html_url: string; pull_request?: unknown; labels: { name: string }[] }[];
103 const newIssues = issues.filter((i) => !i.pull_request).slice(0, 10);
104
105 let processed = 0;
106 let skipped = 0;
107
108 for (const issue of newIssues) {
109 if (await hasFreshDraft(env.CURATED_KV, "triage", String(issue.number), issue.updated_at)) {
110 skipped++;
111 continue;
112 }
113
114 const payload = {
115 number: issue.number,
116 title: issue.title,
117 body: (issue.body ?? "").slice(0, 3000),
118 labels: issue.labels.map((l) => l.name),
119 url: issue.html_url,
120 };
121
122 try {
123 const { content, usage } = await agentChat(
124 [{ role: "system", content: TRIAGE_PROMPT }, { role: "user", content: JSON.stringify(payload) }],
125 env.DEEPSEEK_API_KEY!,
126 true,
127 dsEnv(env)
128 );
129 const parsed = JSON.parse(content) as { bodyEn: string; bodyZh: string };
130 const draft: AgentDraft = {
131 id: String(issue.number),
132 type: "triage",
133 targetNumber: issue.number,
134 targetUrl: issue.html_url,
135 bodyEn: parsed.bodyEn,
136 bodyZh: parsed.bodyZh,
137 generatedAt: new Date().toISOString(),
138 posted: false,
139 };
140 await logUsage(env.CURATED_KV, usage.input, usage.output);
141 if (await saveDraft(env.CURATED_KV, draft, issue.updated_at)) processed++;
142 else skipped++;
143 } catch {
144 skipped++;
145 }
146 }
147
148 return { ok: true, processed, skipped };
149 } catch (e) {
150 return { ok: false, error: String(e) };
151 }
152 }
153
154 async function generatePrReview(env: AgentEnv): Promise<Record<string, unknown>> {
155 const repo = env.GITHUB_REPO ?? "codewhale-hq/CodeWhale";
156 try {
157 const res = await fetch(
158 `https://api.github.com/repos/${repo}/pulls?state=open&sort=created&direction=desc&per_page=20`,
159 {
160 signal: AbortSignal.timeout(OUTBOUND_TIMEOUT_MS),
161 headers: {
162 Accept: "application/vnd.github+json",
163 "X-GitHub-Api-Version": "2022-11-28",
164 "User-Agent": "codewhale-web",
165 ...(env.GITHUB_TOKEN ? { Authorization: `Bearer ${env.GITHUB_TOKEN}` } : {}),
166 },
167 }
168 );
169 const prs = (await res.json()) as { number: number; title: string; body?: string; updated_at: string; html_url: string; changed_files?: number; additions?: number; deletions?: number; user: { login: string } }[];
170
171 let processed = 0;
172 let skipped = 0;
173
174 for (const pr of prs.slice(0, 10)) {
175 if (await hasFreshDraft(env.CURATED_KV, "pr-review", String(pr.number), pr.updated_at)) {
176 skipped++;
177 continue;
178 }
179
180 // Fetch diff stats if not included
181 let diffStats = { changed_files: pr.changed_files ?? 0, additions: pr.additions ?? 0, deletions: pr.deletions ?? 0 };
182 if (!pr.changed_files) {
183 try {
184 const diffRes = await fetch(`https://api.github.com/repos/${repo}/pulls/${pr.number}`, {
185 signal: AbortSignal.timeout(OUTBOUND_TIMEOUT_MS),
186 headers: {
187 Accept: "application/vnd.github+json",
188 "X-GitHub-Api-Version": "2022-11-28",
189 "User-Agent": "codewhale-web",
190 ...(env.GITHUB_TOKEN ? { Authorization: `Bearer ${env.GITHUB_TOKEN}` } : {}),
191 },
192 });
193 const diffData = (await diffRes.json()) as { changed_files?: number; additions?: number; deletions?: number };
194 diffStats = { changed_files: diffData.changed_files ?? 0, additions: diffData.additions ?? 0, deletions: diffData.deletions ?? 0 };
195 } catch { /* use defaults */ }
196 }
197
198 const payload = {
199 number: pr.number,
200 title: pr.title,
201 body: (pr.body ?? "").slice(0, 3000),
202 author: pr.user.login,
203 url: pr.html_url,
204 ...diffStats,
205 };
206
207 try {
208 const { content, usage } = await agentChat(
209 [{ role: "system", content: PR_REVIEW_PROMPT }, { role: "user", content: JSON.stringify(payload) }],
210 env.DEEPSEEK_API_KEY!,
211 true,
212 dsEnv(env)
213 );
214 const parsed = JSON.parse(content) as { bodyEn: string; bodyZh: string };
215 const draft: AgentDraft = {
216 id: String(pr.number),
217 type: "pr-review",
218 targetNumber: pr.number,
219 targetUrl: pr.html_url,
220 bodyEn: parsed.bodyEn,
221 bodyZh: parsed.bodyZh,
222 generatedAt: new Date().toISOString(),
223 posted: false,
224 };
225 await logUsage(env.CURATED_KV, usage.input, usage.output);
226 if (await saveDraft(env.CURATED_KV, draft, pr.updated_at)) processed++;
227 else skipped++;
228 } catch {
229 skipped++;
230 }
231 }
232
233 return { ok: true, processed, skipped };
234 } catch (e) {
235 return { ok: false, error: String(e) };
236 }
237 }
238
239 async function generateStale(env: AgentEnv): Promise<Record<string, unknown>> {
240 const repo = env.GITHUB_REPO ?? "codewhale-hq/CodeWhale";
241 const thirtyDaysAgo = new Date(Date.now() - 30 * 24 * 60 * 60 * 1000).toISOString().slice(0, 10);
242 try {
243 const res = await fetch(
244 `https://api.github.com/search/issues?q=${encodeURIComponent(`repo:${repo} is:issue is:open updated:<${thirtyDaysAgo}`)}&sort=updated&per_page=20`,
245 {
246 signal: AbortSignal.timeout(OUTBOUND_TIMEOUT_MS),
247 headers: {
248 Accept: "application/vnd.github+json",
249 "X-GitHub-Api-Version": "2022-11-28",
250 "User-Agent": "codewhale-web",
251 ...(env.GITHUB_TOKEN ? { Authorization: `Bearer ${env.GITHUB_TOKEN}` } : {}),
252 },
253 }
254 );
255 const data = (await res.json()) as { items?: { number: number; title: string; body?: string; updated_at: string; html_url: string }[] };
256 const issues = data.items ?? [];
257
258 let processed = 0;
259 let skipped = 0;
260
261 for (const issue of issues.slice(0, 10)) {
262 if (await hasFreshDraft(env.CURATED_KV, "stale", String(issue.number), issue.updated_at)) {
263 skipped++;
264 continue;
265 }
266
267 const payload = {
268 number: issue.number,
269 title: issue.title,
270 body: (issue.body ?? "").slice(0, 2000),
271 url: issue.html_url,
272 lastUpdated: issue.updated_at,
273 };
274
275 try {
276 const { content, usage } = await agentChat(
277 [{ role: "system", content: STALE_PROMPT }, { role: "user", content: JSON.stringify(payload) }],
278 env.DEEPSEEK_API_KEY!,
279 true,
280 dsEnv(env)
281 );
282 const parsed = JSON.parse(content) as { bodyEn: string; bodyZh: string };
283 const draft: AgentDraft = {
284 id: String(issue.number),
285 type: "stale",
286 targetNumber: issue.number,
287 targetUrl: issue.html_url,
288 bodyEn: parsed.bodyEn,
289 bodyZh: parsed.bodyZh,
290 generatedAt: new Date().toISOString(),
291 posted: false,
292 };
293 await logUsage(env.CURATED_KV, usage.input, usage.output);
294 if (await saveDraft(env.CURATED_KV, draft, issue.updated_at)) processed++;
295 else skipped++;
296 } catch {
297 skipped++;
298 }
299 }
300
301 return { ok: true, processed, skipped };
302 } catch (e) {
303 return { ok: false, error: String(e) };
304 }
305 }
306
307 async function generateDupes(env: AgentEnv): Promise<Record<string, unknown>> {
308 const repo = env.GITHUB_REPO ?? "codewhale-hq/CodeWhale";
309 try {
310 const res = await fetch(
311 `https://api.github.com/repos/${repo}/issues?state=open&per_page=100`,
312 {
313 signal: AbortSignal.timeout(OUTBOUND_TIMEOUT_MS),
314 headers: {
315 Accept: "application/vnd.github+json",
316 "X-GitHub-Api-Version": "2022-11-28",
317 "User-Agent": "codewhale-web",
318 ...(env.GITHUB_TOKEN ? { Authorization: `Bearer ${env.GITHUB_TOKEN}` } : {}),
319 },
320 }
321 );
322 const issues = (await res.json()) as { number: number; title: string; body?: string; updated_at: string; html_url: string; pull_request?: unknown }[];
323 const openIssues = issues
324 .filter((i) => !i.pull_request)
325 .map((i) => ({
326 number: i.number,
327 title: i.title,
328 body: (i.body ?? "").slice(0, 500),
329 url: i.html_url,
330 }));
331
332 if (openIssues.length < 3) {
333 return { ok: true, skipped: true, reason: "too few issues to compare" };
334 }
335
336 const { content, usage } = await agentChat(
337 [{ role: "system", content: DUPES_PROMPT }, { role: "user", content: JSON.stringify({ issues: openIssues }) }],
338 env.DEEPSEEK_API_KEY!,
339 true,
340 dsEnv(env)
341 );
342
343 const parsed = JSON.parse(content) as { suggestions?: { targetNumber: number; duplicateNumber: number; reason: string; bodyEn: string; bodyZh: string }[] };
344 const suggestions = parsed.suggestions ?? [];
345
346 let processed = 0;
347 for (const s of suggestions) {
348 const draft: AgentDraft = {
349 id: String(s.duplicateNumber),
350 type: "dupes",
351 targetNumber: s.duplicateNumber,
352 bodyEn: s.bodyEn,
353 bodyZh: s.bodyZh,
354 generatedAt: new Date().toISOString(),
355 posted: false,
356 };
357 // saveDraft refuses identities the maintainer already posted or discarded.
358 if (await saveDraft(env.CURATED_KV, draft)) processed++;
359 }
360
361 await logUsage(env.CURATED_KV, usage.input, usage.output);
362 return { ok: true, processed };
363 } catch (e) {
364 return { ok: false, error: String(e) };
365 }
366 }
367
368 async function generateDigest(env: AgentEnv): Promise<Record<string, unknown>> {
369 const repo = env.GITHUB_REPO ?? "codewhale-hq/CodeWhale";
370 const weekAgo = new Date(Date.now() - 7 * 24 * 60 * 60 * 1000).toISOString();
371
372 // Compute week ID
373 const now = new Date();
374 const startOfYear = new Date(now.getFullYear(), 0, 1);
375 const weekNum = Math.ceil(((now.getTime() - startOfYear.getTime()) / 86400000 + startOfYear.getDay() + 1) / 7);
376 const weekId = `${now.getFullYear()}-W${String(weekNum).padStart(2, "0")}`;
377
378 try {
379 if (await getDraftResolution(env.CURATED_KV, "digest", weekId)) {
380 return { ok: true, skipped: true, weekId, reason: "digest already reviewed" };
381 }
382
383 const [issuesRes, pullsRes, stats] = await Promise.all([
384 fetch(
385 `https://api.github.com/repos/${repo}/issues?state=all&since=${weekAgo}&per_page=50&sort=updated&direction=desc`,
386 {
387 signal: AbortSignal.timeout(OUTBOUND_TIMEOUT_MS),
388 headers: {
389 Accept: "application/vnd.github+json",
390 "X-GitHub-Api-Version": "2022-11-28",
391 "User-Agent": "codewhale-web",
392 ...(env.GITHUB_TOKEN ? { Authorization: `Bearer ${env.GITHUB_TOKEN}` } : {}),
393 },
394 }
395 ),
396 fetch(
397 `https://api.github.com/repos/${repo}/pulls?state=all&sort=updated&direction=desc&per_page=50`,
398 {
399 signal: AbortSignal.timeout(OUTBOUND_TIMEOUT_MS),
400 headers: {
401 Accept: "application/vnd.github+json",
402 "X-GitHub-Api-Version": "2022-11-28",
403 "User-Agent": "codewhale-web",
404 ...(env.GITHUB_TOKEN ? { Authorization: `Bearer ${env.GITHUB_TOKEN}` } : {}),
405 },
406 }
407 ),
408 fetchRepoStats(env.GITHUB_TOKEN),
409 ]);
410
411 const issues = (await issuesRes.json()) as { number: number; title: string; state: string; pull_request?: unknown; created_at: string; user: { login: string } }[];
412 const pulls = (await pullsRes.json()) as { number: number; title: string; state: string; merged_at?: string; created_at: string; user: { login: string } }[];
413
414 const weekIssues = issues.filter((i) => !i.pull_request && new Date(i.created_at) > new Date(weekAgo));
415 const weekPRs = pulls.filter((p) => new Date(p.created_at) > new Date(weekAgo));
416 const mergedPRs = pulls.filter((p) => p.merged_at && new Date(p.merged_at) > new Date(weekAgo));
417
418 const contributors = new Set([
419 ...weekIssues.map((i) => i.user.login),
420 ...weekPRs.map((p) => p.user.login),
421 ]);
422
423 const payload = {
424 period: `${weekAgo.slice(0, 10)} — ${new Date().toISOString().slice(0, 10)}`,
425 stats: { stars: stats.stars, forks: stats.forks },
426 newIssues: weekIssues.map((i) => ({ number: i.number, title: i.title, author: i.user.login })),
427 newPRs: weekPRs.map((p) => ({ number: p.number, title: p.title, author: p.user.login })),
428 mergedPRs: mergedPRs.map((p) => ({ number: p.number, title: p.title })),
429 contributors: [...contributors],
430 };
431
432 const { content, usage } = await agentChat(
433 [{ role: "system", content: DIGEST_PROMPT }, { role: "user", content: JSON.stringify(payload) }],
434 env.DEEPSEEK_API_KEY!,
435 true,
436 dsEnv(env)
437 );
438
439 const parsed = JSON.parse(content) as { titleEn: string; titleZh: string; summaryEn: string; summaryZh: string; sections: { heading: string; items: string[] }[] };
440
441 const draft: AgentDraft = {
442 id: weekId,
443 type: "digest",
444 bodyEn: renderDigestBody(parsed, "en"),
445 bodyZh: renderDigestBody(parsed, "zh"),
446 generatedAt: new Date().toISOString(),
447 posted: false,
448 };
449
450 if (!(await saveDraft(env.CURATED_KV, draft))) {
451 return { ok: true, skipped: true, weekId, reason: "digest already reviewed" };
452 }
453
454 // Stage the structured digest for the weekly page, unapproved. The page
455 // renders it only after the maintainer posts the draft from /admin.
456 // Pick fields explicitly: model output must never be able to set
457 // `approved` or any other record key.
458 const record: WeeklyDigestRecord = {
459 titleEn: parsed.titleEn,
460 titleZh: parsed.titleZh,
461 summaryEn: parsed.summaryEn,
462 summaryZh: parsed.summaryZh,
463 sections: parsed.sections,
464 weekId,
465 generatedAt: draft.generatedAt,
466 approved: false,
467 };
468 await env.CURATED_KV?.put(digestRecordKey(weekId), JSON.stringify(record), {
469 expirationTtl: 60 * 60 * 24 * 90,
470 });
471
472 await logUsage(env.CURATED_KV, usage.input, usage.output);
473 return { ok: true, weekId };
474 } catch (e) {
475 return { ok: false, error: String(e) };
476 }
477 }
478
479 /** The existing durable claim authority also serializes each bounded cron batch.
480 * Missing bindings stop generation before provider spend; KV is never a lock.
481 * The 45-minute lease covers at most ten 180-second model calls and reads.
482 * A completed batch retains a short hold for its KV drafts to propagate.
483 */
484 async function withGenerationClaim(env: AgentEnv, task: string, generate: () => Promise<Record<string, unknown>>): Promise<Record<string, unknown>> {
485 if (!env.DEEPSEEK_API_KEY) return { skipped: true, reason: "DEEPSEEK_API_KEY not set" };
486 if (!env.DRAFT_CLAIM_LOCK || !env.CURATED_KV) return { skipped: true, reason: "durable generation claim unavailable; no provider call started" };
487 let lock: DraftClaimLockStub;
488 const token = `generate:${crypto.randomUUID()}`;
489 try {
490 lock = env.DRAFT_CLAIM_LOCK.get(env.DRAFT_CLAIM_LOCK.idFromName(`draft-generation:${task}`));
491 const held = await lock.act({ op: "claim", token, action: "generate", leaseMs: 45 * 60 * 1000 });
492 if (!held.ok) return { ok: true, processed: 0, skipped: 1, reason: "generation already in progress or recently completed" };
493 } catch {
494 return { skipped: true, reason: "durable generation claim unavailable; no provider call started" };
495 }
496 let completed = false;
497 try {
498 const result = await generate();
499 completed = result.ok === true;
500 return result;
501 } finally {
502 // A release failure leaves the durable lease in place until its expiry.
503 await lock.act({ op: "release", token, holdMs: completed ? 120_000 : 0 }).catch(() => undefined);
504 }
505 }
506
507 export const runCurate = (env: AgentEnv) => withGenerationClaim(env, "curate", () => generateCurate(env));
508
509 export const runTriage = (env: AgentEnv) => withGenerationClaim(env, "triage", () => generateTriage(env));
510
511 export const runPrReview = (env: AgentEnv) => withGenerationClaim(env, "prreview", () => generatePrReview(env));
512
513 export const runStale = (env: AgentEnv) => withGenerationClaim(env, "stale", () => generateStale(env));
514
515 export const runDupes = (env: AgentEnv) => withGenerationClaim(env, "dupes", () => generateDupes(env));
516
517 export const runDigest = (env: AgentEnv) => withGenerationClaim(env, "digest", () => generateDigest(env));
518
518 lines TYPESCRIPT