diff --git a/CONTEXT.md b/CONTEXT.md index 6dd32e6..c7cf507 100644 --- a/CONTEXT.md +++ b/CONTEXT.md @@ -77,11 +77,11 @@ A thematic label used to filter and rank Radar Signals, such as Agents and Autom _Avoid_: Signal type, source **Collection Cycle**: -The two-hour interval in which Source Connectors collect new Source Evidence. +One of the two daily collection windows at 09:00 or 17:00 China Standard Time in which Source Connectors collect new Source Evidence before the corresponding Brief Snapshot. _Avoid_: Daily Brief, Observation Window **Publication Time**: -The daily 09:00 China Standard Time when AI Radar automatically publishes its Daily Brief. +The daily 09:00 and 17:00 China Standard Time windows when AI Radar automatically publishes the morning and afternoon Brief Snapshots. _Avoid_: Collection Cycle, user notification time **Correction**: @@ -161,7 +161,7 @@ The publicly accessible, default-configured deployment maintained as the canonic _Avoid_: Self-Hosted Instance, SaaS account **Radar Profile**: -The single collection configuration for a Self-Hosted Instance. It specifies enabled Source Connectors, topic inclusion and exclusion rules, Topic Tags, timing, and output language. +The single collection configuration for a Self-Hosted Instance. It specifies enabled Source Connectors, topic inclusion and exclusion rules, Topic Tags, and output language; the two daily collection windows are fixed product behavior. _Avoid_: User account, Daily Brief **Hosted Multi-Tenant Platform**: @@ -253,7 +253,7 @@ The person holding the deployment-supplied administrator credential that authori _Avoid_: Public visitor, application user **Profile Configuration**: -The non-sensitive, persisted settings of a Radar Profile, including Connector choices, Candidate Filters, scheduling, and language. It is editable in the Profile Configuration Console and exportable or importable. +The non-sensitive, persisted settings of a Radar Profile, including Connector choices, Candidate Filters, and language. The two daily collection and publication windows are fixed product behavior; the remaining settings are editable in the Profile Configuration Console and exportable or importable. _Avoid_: Source Credential, Model Runtime credential **Configuration Version**: diff --git a/README.md b/README.md new file mode 100644 index 0000000..9091e76 --- /dev/null +++ b/README.md @@ -0,0 +1,185 @@ +# Razer-Raders + +Razer-Raders 是面向独立开发者和 AI 产品 Builder 的自托管 AI Radar:从公开开发者来源中发现值得关注的 AI 发展,保留可追溯的 Source Evidence,经过后台评估后生成中文 Daily Brief,帮助你决定一个信号应该「试用、学习、跟进,还是跳过」。 + +项目免费、开源,不依赖中心化 SaaS。每个 Self-Hosted Instance 使用自己的数据库、来源配置和模型运行时;默认不向中心服务发送使用数据、采集证据、Prompt 或模型输出。 + +## 能做什么 + +- 从 GitHub Trending、Hugging Face Trending、Hacker News Show HN 采集公开内容。 +- 在滚动七天 Observation Window 内合并、去重、筛选并排名 Radar Signal。 +- 对候选信号执行证据优先的后台 Assessment Workflow,保留来源、摘要和评估溯源。 +- 每天 09:00 和 17:00(Asia/Shanghai)发布 Brief Snapshot;评估延迟会透明展示,不会静默替换模型或改写已发布 Brief。 +- 提供 Public Brief、Radar Archive、主题筛选、优先级筛选、收藏和深色/浅色模式。 +- 提供响应式 Mobile Reading Experience,支持通过 URL 保留 Archive View State。 +- 提供 Profile Configuration Console,可配置来源连接器、包含/排除词、评估并发和模型运行时。 +- 支持 Ollama 本地模型,以及部署者自己提供的 OpenAI-compatible Chat Completions API。 +- 通过 PostgreSQL 保存 Radar Archive、证据摘要、配置版本、评估任务和 Brief Provenance。 + +## 工作流 + +```text +Source Connectors + ↓ +Candidates → Candidate Filter → Evidence Enrichment + ↓ + Task Worker / Model Runtime + ↓ + Publication Validation / Ranking + ↓ + Immutable Brief Snapshot + ↓ + Web Brief / Radar Archive +``` + +外部页面、仓库、模型卡和社区讨论都被当作 Untrusted Evidence。采集和证据补充统一经过受限 Fetch Gateway,内容不能向系统发出指令,也不能访问凭据或扩大抓取范围。 + +## 快速开始 + +### 前置条件 + +- Docker Desktop 或支持 Docker Compose 的 Docker 环境 +- Node.js 22+、Corepack 和 pnpm 11(仅在本地开发或直接运行脚本时需要) +- 一个模型运行时:本机 Ollama,或一个可访问的 HTTPS OpenAI-compatible API + +### 使用 Docker Compose 启动 + +先在项目根目录准备环境变量。下面的示例使用 Docker Desktop 宿主机上的 Ollama: + +```bash +export POSTGRES_PASSWORD='change-this-password' +export RADAR_ADMIN_TOKEN='change-this-admin-token' +export RADAR_MODEL_RUNTIME='ollama' +export RADAR_OLLAMA_BASE_URL='http://host.docker.internal:11434' +export RADAR_OLLAMA_MODEL='qwen3:8b' + +docker compose up -d --build +``` + +启动流程会依次准备 PostgreSQL、执行数据库迁移、启动 Web Service 和 Task Worker。打开 查看 Public Brief;配置后台位于页面中的「配置后台」入口。 + +如果使用外部 OpenAI-compatible API,将模型运行时替换为: + +```bash +export RADAR_MODEL_RUNTIME='compatible' +export RADAR_COMPATIBLE_RUNTIME_BASE_URL='https://your-provider.example/v1' +export RADAR_COMPATIBLE_RUNTIME_MODEL='your-model' +export RADAR_COMPATIBLE_RUNTIME_API_KEY='your-api-key' + +docker compose up -d --build +``` + +Compatible API 的密钥只从部署环境读取,不会写入 Radar Profile,也不会发送给模型服务以外的地方。生产环境请使用强密码和强管理员 Token,不要继续使用示例值。 + +常用运维命令: + +```bash +docker compose ps +docker compose logs -f web worker +docker compose down +``` + +### 本地开发 + +```bash +pnpm install + +export DATABASE_URL='postgresql://razer_raders:local-development-only@127.0.0.1:5432/razer_raders' +export RADAR_ADMIN_TOKEN='local-development-admin-token' + +docker compose up -d postgres +pnpm db:migrate +pnpm dev +``` + +另开一个终端,在同样的环境变量下启动 Worker: + +```bash +pnpm worker +``` + +开发服务器默认运行在 。如果只需要执行一次采集和评估周期,可以使用: + +```bash +RADAR_WORKER_ONCE=true pnpm worker +``` + +## 配置说明 + +首次启动后,在「配置后台」输入 `RADAR_ADMIN_TOKEN`,加载并配置当前 Radar Profile。可配置内容包括: + +- 启用的 Source Connector +- Candidate 的包含词和排除词 +- 每轮评估上限、模型并发和周期预算 +- Ollama 或 Compatible API 的地址和模型 +- 真实连接测试、Ollama 模型发现、立即采集和延迟 Candidate 重试 + +配置保存为新的不可变版本,并从下一次 Collection Cycle 起生效。手动采集只更新候选和评估队列,不改写当前已发布 Brief。 + +### 主要环境变量 + +| 变量 | 作用 | +| --- | --- | +| `DATABASE_URL` | Web、迁移和 Worker 连接 PostgreSQL 的连接串 | +| `POSTGRES_PASSWORD` | Docker Compose 创建 PostgreSQL 的密码 | +| `POSTGRES_PORT` | Docker Compose 暴露 PostgreSQL 的宿主机端口,默认 `5432` | +| `RADAR_ADMIN_TOKEN` | 配置后台的 Bearer Token;未配置时写操作会禁用 | +| `RADAR_MODEL_RUNTIME` | `ollama` 或 `compatible`,默认 `compatible` | +| `RADAR_OLLAMA_BASE_URL` | Ollama 服务地址 | +| `RADAR_OLLAMA_MODEL` | Ollama 模型名称 | +| `RADAR_COMPATIBLE_RUNTIME_BASE_URL` | OpenAI-compatible API 地址 | +| `RADAR_COMPATIBLE_RUNTIME_MODEL` | 外部 API 使用的模型名称 | +| `RADAR_COMPATIBLE_RUNTIME_API_KEY` | 外部 API 密钥,仅在服务端环境使用 | +| `RADAR_INCLUDE_TERMS` / `RADAR_EXCLUDE_TERMS` | 初始化 Profile 时使用的逗号分隔筛选词 | +| `RADAR_PIPELINE_VERSION` | Assessment Pipeline 版本标识,用于发布溯源 | + +## 测试与检查 + +```bash +pnpm lint +pnpm test +pnpm build +``` + +端到端测试会自动创建独立的 PostgreSQL 容器,并验证 Brief Publication 与浏览器流程: + +```bash +pnpm test:e2e +``` + +还可以运行浏览器测试或 HTTPS 反向代理测试: + +```bash +pnpm test:e2e:browser +pnpm test:https-edge +``` + +## 项目结构 + +```text +src/ +├── app/ Next.js 页面与 API Route +├── components/ Brief、Archive、配置后台和 UI 组件 +├── db/migrations/ PostgreSQL 迁移 +├── lib/radar/ 采集、筛选、证据、评估、发布和检索领域逻辑 +└── worker.ts 定时采集、评估和 Brief 发布 Worker +test/ 单元、集成、E2E 和浏览器测试 +compose*.yaml 本地、生产式和 HTTPS Compose 配置 +scripts/ E2E 与 HTTPS 边缘测试脚本 +docs/adr/ 架构决策记录 +``` + +技术栈为 Next.js、React、TypeScript、Node.js、PostgreSQL 和 Docker Compose。 + +## 当前范围 + +- 当前产品是单一 Self-Hosted Instance 和单一 Radar Profile,不是 Hosted Multi-Tenant Platform。 +- MVP 仅包含 GitHub Trending、Hugging Face Trending 和 Show HN 三个 Source Connector。 +- Public Brief 展示已发布的 Brief Snapshot,不是未经评估的实时原始信息流。 +- Grounded Assessment 依赖可用且正确配置的模型运行时;运行时不可用时会显示评估延迟。 +- Source Connector 受外部站点访问规则、网络状况和 rate limit 影响,Connector Health 会在 Brief 中公开摘要。 +- 第三方内容和链接仍受其各自许可证及服务条款约束;本项目只保留必要的证据摘要,不存储完整第三方作品。 + +## 许可证 + +本项目使用 [Apache License 2.0](LICENSE)。 diff --git a/compose.yaml b/compose.yaml index 9f035a5..f5c909b 100644 --- a/compose.yaml +++ b/compose.yaml @@ -62,7 +62,6 @@ services: condition: service_completed_successfully environment: DATABASE_URL: postgresql://razer_raders:${POSTGRES_PASSWORD:-local-development-only}@postgres:5432/razer_raders - RADAR_COLLECTION_INTERVAL_MS: 7200000 RADAR_CONFIGURATION_VERSION: ${RADAR_CONFIGURATION_VERSION:-profile@v1} RADAR_MODEL_RUNTIME: ${RADAR_MODEL_RUNTIME:-compatible} RADAR_COMPATIBLE_RUNTIME_API_KEY: ${RADAR_COMPATIBLE_RUNTIME_API_KEY:-} diff --git a/src/app/api/profile/retry-delayed/route.ts b/src/app/api/profile/retry-delayed/route.ts new file mode 100644 index 0000000..0007b0e --- /dev/null +++ b/src/app/api/profile/retry-delayed/route.ts @@ -0,0 +1,13 @@ +import { getRequiredRadarProfile } from "@/lib/radar/profile-archive"; +import { createProfileReassessHandler } from "@/lib/radar/profile-route"; +import { postgresCandidateTaskArchive } from "@/lib/radar/task-queue"; + +export const POST = createProfileReassessHandler({ + requeue: async () => { + const profile = await getRequiredRadarProfile(); + return postgresCandidateTaskArchive.requeueDelayedAssessments?.({ + configurationVersion: profile.id, + runtimeId: `${profile.runtime.kind}:${profile.runtime.model}`, + }) ?? 0; + }, +}); diff --git a/src/components/brief-presentation.ts b/src/components/brief-presentation.ts index 38442d6..3528498 100644 --- a/src/components/brief-presentation.ts +++ b/src/components/brief-presentation.ts @@ -17,6 +17,9 @@ export function getBriefHeading(input: BriefPresentationInput) { } export function getAssessmentBanner(input: BriefPresentationInput) { + if (input.assessmentDelay && input.hasPublishedSignals) { + return `另有 ${input.assessmentDelay.candidateCount} 个 Candidate 评估延迟,当前已发布日报正常展示。原因:${input.assessmentDelay.detail}`; + } if (input.availability === "evaluating" && input.hasPublishedSignals) { return `另有 ${input.pendingCandidateCount} 个新 Candidate 正在评估,不会混入当前已发布日报。`; } diff --git a/src/components/profile-config.tsx b/src/components/profile-config.tsx index c4ee4fe..a64494c 100644 --- a/src/components/profile-config.tsx +++ b/src/components/profile-config.tsx @@ -16,7 +16,6 @@ type ManualCollectionPayload = { connectorResults?: readonly { status: "failed" function configuration(profile: RadarProfileConfig): RadarProfileConfig { return { - collectionIntervalMs: profile.collectionIntervalMs, enabledConnectorIds: profile.enabledConnectorIds, excludeTerms: profile.excludeTerms, includeTerms: profile.includeTerms, @@ -135,6 +134,18 @@ export function ProfileConfig({ connectors }: { connectors: readonly RadarConnec } }; + const retryDelayed = async () => { + setPending(true); + try { + const result = await request("/api/profile/retry-delayed", "POST") as { requeuedCount: number }; + setMessage(result.requeuedCount > 0 ? `已重新排队 ${result.requeuedCount} 个延迟 Candidate,等待 Worker 处理。` : "当前没有可重试的延迟 Candidate。"); + } catch (error) { + setMessage(error instanceof Error ? error.message : "重试延迟 Candidate 失败。"); + } finally { + setPending(false); + } + }; + const isOllama = profile?.runtime.kind === "ollama"; return
@@ -143,7 +154,7 @@ export function ProfileConfig({ connectors }: { connectors: readonly RadarConnec
00管理员认证
setToken(event.target.value)} placeholder="RADAR_ADMIN_TOKEN" type="password" value={token} />

Token 仅保留在当前浏览器会话,不写入 Profile 或发送给模型服务。

{profile ? <>
01来源连接器
{CONNECTORS.map(({ id, label }) => )}
-
02筛选与采集

立即采集只更新候选与评估队列,不发布或改写当天 Daily Brief。

+
02筛选与采集

Worker 每天 09:00 和 17:00 自动采集并发布;手动采集只更新队列,不改写当前 Brief。重试按钮会重新排队最近 7 天内的评估延迟 Candidate。

03模型运行时
{isOllama ? : null}

{isOllama ? "发现模型会更新上方的 Ollama 模型列表;连接测试不会保存配置。" : "Compatible API Key 只读取部署环境;切换运行时不会覆盖当前地址或模型,请填写对应服务的 HTTPS 地址与模型后再测试。"}

: null} diff --git a/src/components/radar-app.tsx b/src/components/radar-app.tsx index 75f79c6..f766bc6 100644 --- a/src/components/radar-app.tsx +++ b/src/components/radar-app.tsx @@ -192,7 +192,7 @@ export function RadarApp({ brief }: { brief: RadarBrief }) {
系统就绪
- Asia/Shanghai · 09:00 发布 + Asia/Shanghai · 09:00 / 17:00 发布
diff --git a/src/db/migrations/015_twice_daily_briefs.sql b/src/db/migrations/015_twice_daily_briefs.sql new file mode 100644 index 0000000..7e11aa4 --- /dev/null +++ b/src/db/migrations/015_twice_daily_briefs.sql @@ -0,0 +1,18 @@ +ALTER TABLE brief_snapshots + ADD COLUMN IF NOT EXISTS publication_slot TEXT NOT NULL DEFAULT 'morning'; + +UPDATE brief_snapshots +SET publication_slot = 'morning'; + +ALTER TABLE brief_snapshots + DROP CONSTRAINT IF EXISTS brief_snapshots_publication_slot_check; + +ALTER TABLE brief_snapshots + ADD CONSTRAINT brief_snapshots_publication_slot_check + CHECK (publication_slot IN ('morning', 'afternoon')); + +DROP INDEX IF EXISTS brief_snapshots_published_day_idx; + +CREATE UNIQUE INDEX IF NOT EXISTS brief_snapshots_published_day_slot_idx + ON brief_snapshots (publication_day, publication_slot) + WHERE status = 'published'; diff --git a/src/db/migrations/016_normalize_legacy_brief_slots.sql b/src/db/migrations/016_normalize_legacy_brief_slots.sql new file mode 100644 index 0000000..5783fa6 --- /dev/null +++ b/src/db/migrations/016_normalize_legacy_brief_slots.sql @@ -0,0 +1,16 @@ +DROP INDEX IF EXISTS brief_snapshots_published_day_slot_idx; + +WITH ranked_snapshots AS ( + SELECT id, + ROW_NUMBER() OVER (PARTITION BY publication_day ORDER BY published_at, id) AS slot_rank + FROM brief_snapshots + WHERE status = 'published' +) +UPDATE brief_snapshots snapshot +SET publication_slot = CASE WHEN ranked.slot_rank = 1 THEN 'morning' ELSE 'afternoon' END +FROM ranked_snapshots ranked +WHERE snapshot.id = ranked.id; + +CREATE UNIQUE INDEX IF NOT EXISTS brief_snapshots_published_day_slot_idx + ON brief_snapshots (publication_day, publication_slot) + WHERE status = 'published'; diff --git a/src/lib/radar/archive.ts b/src/lib/radar/archive.ts index 071c897..8270514 100644 --- a/src/lib/radar/archive.ts +++ b/src/lib/radar/archive.ts @@ -57,7 +57,7 @@ export async function getLatestPublishedBrief(): Promise `SELECT id, published_at, configuration_version, ranking_policy_version, model_runtime_id, pipeline_version FROM brief_snapshots WHERE status = 'published' - ORDER BY publication_day DESC + ORDER BY published_at DESC LIMIT 1`, ); const snapshot = brief.rows[0]; diff --git a/src/lib/radar/brief-contract.ts b/src/lib/radar/brief-contract.ts index 6963c38..e1c45c5 100644 --- a/src/lib/radar/brief-contract.ts +++ b/src/lib/radar/brief-contract.ts @@ -79,8 +79,16 @@ export function createArchiveRadarBrief(input: { ? { candidateCount: assessment.candidateCount, detail: assessment.detail } : undefined; + const availability = assessment.status === "assessment-delayed" && !brief + ? "assessment-delayed" + : assessment.status === "evaluating" + ? "evaluating" + : brief + ? "published" + : "unpublished"; + return { - availability: assessment.status === "assessment-delayed" ? "assessment-delayed" : assessment.status === "evaluating" ? "evaluating" : brief ? "published" : "unpublished", + availability, ...(assessmentDelay ? { assessmentDelay } : {}), connectors, ...(brief?.coverage ? { coverage: brief.coverage } : {}), diff --git a/src/lib/radar/brief-publication-archive.ts b/src/lib/radar/brief-publication-archive.ts index 49fd225..8ce0e71 100644 --- a/src/lib/radar/brief-publication-archive.ts +++ b/src/lib/radar/brief-publication-archive.ts @@ -2,6 +2,7 @@ import type { QueryResultRow } from "pg"; import type { PublicationArchive, PublicationCandidate, PublishedSignalInput, ReadyPublicationAssessment } from "./brief-publication.ts"; import type { EvidenceFirstAssessment } from "./assessment-contract.ts"; import type { BriefProvenance } from "./brief-contract.ts"; +import type { PublicationSlot } from "./daily-publication-schedule.ts"; import { getDatabasePool, withTransaction } from "./database.ts"; type CandidateRow = QueryResultRow & { @@ -131,10 +132,10 @@ export const postgresBriefPublicationArchive: PublicationArchive = { } satisfies ReadyPublicationAssessment)); }, - async hasPublishedBrief(publicationDay) { + async hasPublishedBrief(publicationDay, publicationSlot = "morning" as PublicationSlot) { const result = await getDatabasePool().query( - "SELECT 1 FROM brief_snapshots WHERE status = 'published' AND publication_day = $1 LIMIT 1", - [publicationDay], + "SELECT 1 FROM brief_snapshots WHERE status = 'published' AND publication_day = $1 AND publication_slot = $2 LIMIT 1", + [publicationDay, publicationSlot], ); return (result.rowCount ?? 0) > 0; }, @@ -148,10 +149,10 @@ export const postgresBriefPublicationArchive: PublicationArchive = { ); }, - async publishBrief({ id, provenance, publicationDay, publishedAt, signals }: { id: string; provenance: BriefProvenance; publicationDay: string; publishedAt: string; signals: readonly PublishedSignalInput[] }) { + async publishBrief({ id, provenance, publicationDay, publicationSlot, publishedAt, signals }: { id: string; provenance: BriefProvenance; publicationDay: string; publicationSlot: PublicationSlot; publishedAt: string; signals: readonly PublishedSignalInput[] }) { return withTransaction(async (client) => { - await client.query("SELECT pg_advisory_xact_lock(hashtext('razer-raders:brief:' || $1))", [publicationDay]); - const existing = await client.query("SELECT 1 FROM brief_snapshots WHERE status = 'published' AND publication_day = $1 LIMIT 1", [publicationDay]); + await client.query("SELECT pg_advisory_xact_lock(hashtext('razer-raders:brief:' || $1 || ':' || $2))", [publicationDay, publicationSlot]); + const existing = await client.query("SELECT 1 FROM brief_snapshots WHERE status = 'published' AND publication_day = $1 AND publication_slot = $2 LIMIT 1", [publicationDay, publicationSlot]); if (existing.rowCount) return "already-published" as const; const coverage = await client.query( @@ -174,12 +175,13 @@ export const postgresBriefPublicationArchive: PublicationArchive = { await client.query( `INSERT INTO brief_snapshots ( - id, published_at, publication_day, status, configuration_version, ranking_policy_version, model_runtime_id, pipeline_version - ) VALUES ($1, $2, $3, 'published', $4, $5, $6, $7)`, + id, published_at, publication_day, publication_slot, status, configuration_version, ranking_policy_version, model_runtime_id, pipeline_version + ) VALUES ($1, $2, $3, $4, 'published', $5, $6, $7, $8)`, [ id, publishedAt, publicationDay, + publicationSlot, provenance.configurationVersion, provenance.rankingPolicyVersion, provenance.modelRuntimeId, diff --git a/src/lib/radar/brief-publication.ts b/src/lib/radar/brief-publication.ts index 429f4db..30bdbe6 100644 --- a/src/lib/radar/brief-publication.ts +++ b/src/lib/radar/brief-publication.ts @@ -1,6 +1,6 @@ import type { AssessmentEvidence, AssessmentWithContent, EvidenceFirstAssessment, GroundedAssessment, ModelRuntime } from "./assessment-contract.ts"; import type { Priority, SignalState } from "../../components/radar-data.ts"; -import { getCstDay } from "./daily-publication-schedule.ts"; +import { getCstDay, type PublicationSlot } from "./daily-publication-schedule.ts"; import { MAX_DAILY_BRIEF_SIGNALS, type BriefProvenance } from "./brief-contract.ts"; export type PublicationCandidate = { @@ -40,6 +40,7 @@ type PublishBriefInput = { id: string; provenance: BriefProvenance; publicationDay: string; + publicationSlot: PublicationSlot; publishedAt: string; signals: readonly PublishedSignalInput[]; }; @@ -47,7 +48,7 @@ type PublishBriefInput = { export type PublicationArchive = { getCandidatesForPublication: (limit?: number) => Promise; getReadyAssessments?: (limit?: number) => Promise; - hasPublishedBrief: (publicationDay: string) => Promise; + hasPublishedBrief: (publicationDay: string, publicationSlot?: PublicationSlot) => Promise; markCandidateAssessmentDelayed: (input: { candidateId: string; detail: string }) => Promise; publishBrief: (input: PublishBriefInput) => Promise<"already-published" | "published">; recordPipelineStage: (input: { collectionRunId?: string; detail?: string; publicationDay: string; stage: PipelineStage; status: PipelineStageStatus }) => Promise; @@ -191,15 +192,17 @@ export function createReadyBriefPublisher(input: { createBriefId: () => string; isCitationAccessible: CitationAccessibility; maxAssessments?: number; + publicationSlot?: PublicationSlot; pipelineVersion: string; }) { const { archive, clock, createBriefId, isCitationAccessible, pipelineVersion } = input; const maxAssessments = Math.min(input.maxAssessments ?? MAX_DAILY_BRIEF_SIGNALS, MAX_DAILY_BRIEF_SIGNALS); + const publicationSlot = input.publicationSlot ?? "morning"; return { async publishDailyBrief(): Promise { const publishedAt = clock(); const publicationDay = getCstDay(publishedAt); - if (await archive.hasPublishedBrief(publicationDay)) return { status: "already-published" }; + if (await archive.hasPublishedBrief(publicationDay, publicationSlot)) return { status: "already-published" }; const ready = rankReadyAssessments(await archive.getReadyAssessments?.(maxAssessments) ?? []); if (!ready.length) return { reason: "Observation Window 内没有已评估待发布的 Candidate。", status: "rejected" }; const first = ready[0]!; @@ -222,6 +225,7 @@ export function createReadyBriefPublisher(input: { id, provenance: { configurationVersion: first.configurationVersion, modelRuntimeId: first.runtimeId, pipelineVersion, rankingPolicyVersion: first.candidate.rankingPolicyVersion }, publicationDay, + publicationSlot, publishedAt: publishedAt.toISOString(), signals, }); @@ -239,6 +243,7 @@ export function createBriefPublisher(input: { createBriefId: () => string; isCitationAccessible: CitationAccessibility; maxAssessments?: number; + publicationSlot?: PublicationSlot; pipelineVersion: string; runtime: ModelRuntime; }) { @@ -251,6 +256,7 @@ export function createBriefPublisher(input: { pipelineVersion, runtime, } = input; + const publicationSlot = input.publicationSlot ?? "morning"; const assessmentBudgetMs = input.assessmentBudgetMs ?? 30 * 60 * 1000; const assessmentConcurrency = input.assessmentConcurrency ?? 1; const maxAssessments = input.maxAssessments ?? 10; @@ -258,7 +264,7 @@ export function createBriefPublisher(input: { async function publishDailyBrief(): Promise { const publishedAt = clock(); const publicationDay = getCstDay(publishedAt); - if (await archive.hasPublishedBrief(publicationDay)) { + if (await archive.hasPublishedBrief(publicationDay, publicationSlot)) { await archive.recordPipelineStage({ detail: "当日 Brief 已发布,跳过重复发布。", publicationDay, stage: "publication", status: "succeeded" }); return { status: "already-published" }; } @@ -356,6 +362,7 @@ export function createBriefPublisher(input: { rankingPolicyVersion: rankingPolicyVersions[0]!, }, publicationDay, + publicationSlot, publishedAt: publishedAt.toISOString(), signals, }); diff --git a/src/lib/radar/daily-publication-schedule.ts b/src/lib/radar/daily-publication-schedule.ts index 65c7402..48d8708 100644 --- a/src/lib/radar/daily-publication-schedule.ts +++ b/src/lib/radar/daily-publication-schedule.ts @@ -1,28 +1,53 @@ const chinaStandardTimeOffsetMs = 8 * 60 * 60 * 1000; const dayMs = 24 * 60 * 60 * 1000; -const publicationHourUtc = 1; + +export type PublicationSlot = "morning" | "afternoon"; +export type DuePublication = { day: string; slot: PublicationSlot }; + +const publicationSlots: readonly { hourUtc: number; slot: PublicationSlot }[] = [ + { hourUtc: 1, slot: "morning" }, + { hourUtc: 9, slot: "afternoon" }, +]; export function getCstDay(value: Date) { return new Date(value.getTime() + chinaStandardTimeOffsetMs).toISOString().slice(0, 10); } -function publicationAtForCstDay(day: string) { - return new Date(`${day}T${String(publicationHourUtc).padStart(2, "0")}:00:00.000Z`); +function publicationAtForCstDay(day: string, hourUtc: number) { + return new Date(`${day}T${String(hourUtc).padStart(2, "0")}:00:00.000Z`); } export function createDailyPublicationSchedule(clock: () => Date) { + const getToday = () => getCstDay(clock()); + const getDuePublications = (): readonly DuePublication[] => { + const now = clock(); + const day = getCstDay(now); + return publicationSlots + .filter(({ hourUtc }) => now >= publicationAtForCstDay(day, hourUtc)) + .map(({ slot }) => ({ day, slot })); + }; + return { + getDuePublications, + + getDuePublication(): DuePublication | null { + return getDuePublications().at(-1) ?? null; + }, + getDuePublicationDay() { - const now = clock(); - const today = getCstDay(now); - return now >= publicationAtForCstDay(today) ? today : null; + return this.getDuePublication()?.day ?? null; }, getNextPublicationAt() { const now = clock(); - const today = getCstDay(now); - const todayPublication = publicationAtForCstDay(today); - return now < todayPublication ? todayPublication : new Date(todayPublication.getTime() + dayMs); + const today = getToday(); + const nextToday = publicationSlots + .map(({ hourUtc }) => publicationAtForCstDay(today, hourUtc)) + .find((publicationAt) => now < publicationAt); + if (nextToday) return nextToday; + + const tomorrow = new Date(publicationAtForCstDay(today, publicationSlots[0]!.hourUtc).getTime() + dayMs); + return tomorrow; }, }; } diff --git a/src/lib/radar/radar-profile.ts b/src/lib/radar/radar-profile.ts index a4e9c5e..d884d50 100644 --- a/src/lib/radar/radar-profile.ts +++ b/src/lib/radar/radar-profile.ts @@ -12,7 +12,6 @@ export type RadarRuntimeConfig = { }; export type RadarProfileConfig = { - collectionIntervalMs: number; enabledConnectorIds: readonly ConnectorId[]; excludeTerms: readonly string[]; includeTerms: readonly string[]; @@ -25,7 +24,6 @@ export type RadarProfile = RadarProfileConfig & { }; export type RadarProfileEnvironment = { - RADAR_COLLECTION_INTERVAL_MS?: string; RADAR_COMPATIBLE_RUNTIME_BASE_URL?: string; RADAR_COMPATIBLE_RUNTIME_MODEL?: string; RADAR_EXCLUDE_TERMS?: string; @@ -98,7 +96,6 @@ export function parseRadarProfileConfig(value: unknown): RadarProfileConfig { throw new Error("至少需要启用一个 Connector。"); } return { - collectionIntervalMs: parseInteger(profile.collectionIntervalMs, "collectionIntervalMs", 60_000, 86_400_000), enabledConnectorIds, excludeTerms: parseStringList(profile.excludeTerms, "excludeTerms"), includeTerms: parseStringList(profile.includeTerms, "includeTerms"), @@ -118,7 +115,6 @@ export function createInitialRadarProfileConfig(environment: RadarProfileEnviron model: environment.RADAR_COMPATIBLE_RUNTIME_MODEL ?? "not-configured", }; return parseRadarProfileConfig({ - collectionIntervalMs: Number(environment.RADAR_COLLECTION_INTERVAL_MS ?? 7_200_000), enabledConnectorIds: ["github-trending", "hugging-face-trending", "show-hn"], excludeTerms: readEnvironmentTerms(environment.RADAR_EXCLUDE_TERMS), includeTerms: readEnvironmentTerms(environment.RADAR_INCLUDE_TERMS), diff --git a/src/lib/radar/task-queue.ts b/src/lib/radar/task-queue.ts index d661ba9..7123a13 100644 --- a/src/lib/radar/task-queue.ts +++ b/src/lib/radar/task-queue.ts @@ -127,13 +127,14 @@ function priorityFor(builderValue?: AssessmentWithContent["builderValue"]): Publ } export type CandidateTaskArchive = { - claimNext: (input: { leaseMs: number; now: Date; preferredKind?: CandidateTaskKind; workerId: string }) => Promise; + claimNext: (input: { excludeTaskIds?: readonly string[]; leaseMs: number; now: Date; preferredKind?: CandidateTaskKind; workerId: string }) => Promise; completeAssessment: (input: { assessment: GroundedAssessment; task: ClaimedCandidateTask }) => Promise; completeEnrichment: (input: { result: EvidenceEnrichmentResult; task: ClaimedCandidateTask }) => Promise; enqueueEnrichment: (input: { candidate: Candidate; configurationVersion: string; force?: boolean; runtimeId: string }) => Promise; fail: (input: { errorMessage: string; task: ClaimedCandidateTask }) => Promise<"delayed" | "retryable">; getStatistics: (input: { cycleStartedAt: Date }) => Promise; release: (input: { task: ClaimedCandidateTask }) => Promise; + requeueDelayedAssessments?: (input: { configurationVersion: string; runtimeId: string }) => Promise; requeueReadyAssessments: (input: { configurationVersion: string; runtimeId: string }) => Promise; }; @@ -171,7 +172,7 @@ export const postgresCandidateTaskArchive: CandidateTaskArchive = { }); }, - async claimNext({ leaseMs, now, preferredKind, workerId }) { + async claimNext({ excludeTaskIds = [], leaseMs, now, preferredKind, workerId }) { const leaseExpiresAt = new Date(now.getTime() + leaseMs); const task = await withTransaction(async (client) => { await client.query( @@ -187,6 +188,7 @@ export const postgresCandidateTaskArchive: CandidateTaskArchive = { JOIN radar_candidates candidate ON candidate.id = task.candidate_id WHERE task.status IN ('queued', 'retryable') AND candidate.last_collected_at >= $1::timestamptz - INTERVAL '7 days' + AND task.id <> ALL($5::text[]) ORDER BY CASE WHEN $4::text IS NOT NULL AND task.task_kind = $4 THEN 0 ELSE 1 END, EXISTS (SELECT 1 FROM candidate_evidence_digests digest_link WHERE digest_link.candidate_id = candidate.id) DESC, (SELECT COUNT(*) FROM candidate_source_evidence source_link WHERE source_link.candidate_id = candidate.id) DESC, @@ -200,7 +202,7 @@ export const postgresCandidateTaskArchive: CandidateTaskArchive = { FROM next_task WHERE task.id = next_task.id RETURNING task.id, task.candidate_id, task.task_kind, task.evidence_fingerprint, task.configuration_version, task.runtime_id, task.attempt_count, task.claimed_by, task.claimed_at`, - [now, leaseExpiresAt, workerId, preferredKind ?? null], + [now, leaseExpiresAt, workerId, preferredKind ?? null, excludeTaskIds], ); const claimed = result.rows[0] ?? null; if (claimed) { @@ -437,6 +439,42 @@ export const postgresCandidateTaskArchive: CandidateTaskArchive = { return candidates.rowCount ?? 0; }); }, + + async requeueDelayedAssessments({ configurationVersion, runtimeId }) { + return withTransaction(async (client) => { + const candidates = await client.query<{ id: string }>( + `SELECT id FROM radar_candidates + WHERE evaluation_status = 'assessment-delayed' + AND last_collected_at >= NOW() - INTERVAL '7 days' + FOR UPDATE`, + ); + for (const candidate of candidates.rows) { + const pendingEnrichment = await client.query( + `UPDATE candidate_tasks + SET status = 'queued', configuration_version = $2, runtime_id = $3, last_error = NULL, + claimed_by = NULL, claimed_at = NULL, lease_expires_at = NULL, completed_at = NULL + WHERE candidate_id = $1 AND task_kind = 'enrichment' AND status IN ('queued', 'retryable') + RETURNING id`, + [candidate.id, configurationVersion, runtimeId], + ); + if (!pendingEnrichment.rowCount) { + await client.query( + `INSERT INTO candidate_tasks (id, candidate_id, task_kind, status, evidence_fingerprint, configuration_version, runtime_id) + VALUES ($1, $2, 'enrichment', 'queued', $3, $4, $5)`, + [randomUUID(), candidate.id, `manual-retry:${randomUUID()}`, configurationVersion, runtimeId], + ); + } + await client.query( + `UPDATE radar_candidates + SET lifecycle_status = '待补证', evaluation_status = 'queued', assessment_delay_detail = NULL, + assessment_result = NULL, assessment_fingerprint = NULL, assessment_task_id = NULL, updated_at = NOW() + WHERE id = $1`, + [candidate.id], + ); + } + return candidates.rowCount ?? 0; + }); + }, }; function failureMessage(error: unknown) { @@ -463,15 +501,15 @@ export function createCandidateTaskWorker(input: { async runCycle() { const deadline = clock().getTime() + timeBudgetMs; let completed = 0; - let runtimeUnavailable = false; const claimedKinds = new Set(); - while (completed < maxTasks && clock().getTime() < deadline && !runtimeUnavailable) { + const skippedTaskIds = new Set(); + while (completed < maxTasks && clock().getTime() < deadline) { const batch: ClaimedCandidateTask[] = []; while (batch.length < concurrency && completed + batch.length < maxTasks && clock().getTime() < deadline) { const preferredKind = claimedKinds.has("assessment") ? claimedKinds.has("enrichment") ? undefined : "enrichment" : claimedKinds.has("enrichment") ? "assessment" : undefined; - const task = await archive.claimNext({ leaseMs, now: clock(), preferredKind, workerId }); + const task = await archive.claimNext({ excludeTaskIds: [...skippedTaskIds], leaseMs, now: clock(), preferredKind, workerId }); if (!task) break; batch.push(task); claimedKinds.add(task.kind); @@ -486,7 +524,7 @@ export function createCandidateTaskWorker(input: { } else { const taskRuntime = input.getRuntime ? await input.getRuntime(task.configurationVersion) : runtime ?? null; if (!taskRuntime || taskRuntime.id !== task.runtimeId) { - runtimeUnavailable = true; + skippedTaskIds.add(task.id); await archive.release({ task }); return; } diff --git a/src/lib/radar/task-worker-schedule.ts b/src/lib/radar/task-worker-schedule.ts index 7b23e96..d148d64 100644 --- a/src/lib/radar/task-worker-schedule.ts +++ b/src/lib/radar/task-worker-schedule.ts @@ -1,24 +1,20 @@ import { createDailyPublicationSchedule } from "./daily-publication-schedule.ts"; type WorkerTimers = { - clearInterval: (handle: TimerHandle) => void; clearTimeout: (handle: TimerHandle) => void; - setInterval: (callback: () => void, delay: number) => TimerHandle; setTimeout: (callback: () => void, delay: number) => TimerHandle; }; export function createTaskWorkerSchedule(input: { clock: () => Date; - collectionIntervalMs: number; collect: () => Promise<"failed" | "succeeded">; - getCollectionIntervalMs?: () => Promise; + getNextCycleAt: () => Date; onError?: (error: unknown) => void; publish: () => Promise; timers: WorkerTimers; }) { - const { clock, collectionIntervalMs, collect, getCollectionIntervalMs, onError = console.error, publish, timers } = input; - let collectionTimer: TimerHandle | undefined; - let publicationTimer: TimerHandle | undefined; + const { clock, collect, getNextCycleAt, onError = console.error, publish, timers } = input; + let cycleTimer: TimerHandle | undefined; let stopped = false; const collectSafely = async () => { @@ -40,7 +36,7 @@ export function createTaskWorkerSchedule(input: { } }; - const collectAtStartup = async () => { + const runCycle = async () => { const collectionStatus = await collectSafely(); if (collectionStatus === "succeeded" && createDailyPublicationSchedule(clock).getDuePublicationDay() && !await publishSafely()) { return "failed" as const; @@ -48,41 +44,26 @@ export function createTaskWorkerSchedule(input: { return collectionStatus; }; - const scheduleNextPublication = () => { + const scheduleNextCycle = () => { if (stopped) return; - const delay = createDailyPublicationSchedule(clock).getNextPublicationAt().getTime() - clock().getTime(); - publicationTimer = timers.setTimeout(() => { - void publishSafely().then(() => { if (!stopped) scheduleNextPublication(); }); - }, delay); - }; - - const scheduleNextCollection = async () => { - let delay = collectionIntervalMs; - try { - if (getCollectionIntervalMs) delay = await getCollectionIntervalMs(); - } catch (error) { - onError(error); - } - if (stopped) return; - collectionTimer = timers.setTimeout(() => { - void collectSafely().then(() => { void scheduleNextCollection(); }); + const delay = Math.max(0, getNextCycleAt().getTime() - clock().getTime()); + cycleTimer = timers.setTimeout(() => { + void runCycle().then(() => { if (!stopped) scheduleNextCycle(); }); }, delay); }; return { - runOnce: collectAtStartup, + runOnce: runCycle, async start() { stopped = false; - await collectAtStartup(); - await scheduleNextCollection(); - scheduleNextPublication(); + await runCycle(); + scheduleNextCycle(); }, stop() { stopped = true; - if (collectionTimer !== undefined) timers.clearTimeout(collectionTimer); - if (publicationTimer !== undefined) timers.clearTimeout(publicationTimer); + if (cycleTimer !== undefined) timers.clearTimeout(cycleTimer); }, }; } diff --git a/src/worker.ts b/src/worker.ts index 85bb4f2..1efba59 100644 --- a/src/worker.ts +++ b/src/worker.ts @@ -9,8 +9,6 @@ import { getDatabasePool } from "./lib/radar/database.ts"; import { getRequiredRadarProfile } from "./lib/radar/profile-archive.ts"; import { createTaskWorkerSchedule } from "./lib/radar/task-worker-schedule.ts"; -const defaultCollectionIntervalMs = 2 * 60 * 60 * 1000; - async function collectScheduledSources() { const result = await collectConfiguredSources(); if (result.status === "already-running") { @@ -20,7 +18,7 @@ async function collectScheduledSources() { return result.status; } -async function publishDailyBriefIfConfigured() { +async function publishDailyBriefIfConfigured(publicationSlot: "morning" | "afternoon") { await getRequiredRadarProfile(); const result = await createReadyBriefPublisher({ archive: postgresBriefPublicationArchive, @@ -28,6 +26,7 @@ async function publishDailyBriefIfConfigured() { createBriefId: randomUUID, isCitationAccessible: (url) => createCitationAccessibilityCheck([url])(url), maxAssessments: MAX_DAILY_BRIEF_SIGNALS, + publicationSlot, pipelineVersion: process.env.RADAR_PIPELINE_VERSION ?? "evidence-first-assessment@v1", }).publishDailyBrief(); if (result.status === "published") console.log(`日报已发布:${result.signalCount} 个信号`); @@ -37,19 +36,24 @@ async function publishDailyBriefIfConfigured() { } async function publishDailyBriefWhenDue() { - if (!createDailyPublicationSchedule(() => new Date()).getDuePublicationDay()) return; - await publishDailyBriefIfConfigured(); + const due = createDailyPublicationSchedule(() => new Date()).getDuePublications(); + for (const publication of due) { + try { + await publishDailyBriefIfConfigured(publication.slot); + } catch (error) { + console.error(`${publication.slot} 时段日报任务失败:`, error); + } + } } async function runWorker() { const schedule = createTaskWorkerSchedule({ clock: () => new Date(), - collectionIntervalMs: defaultCollectionIntervalMs, - getCollectionIntervalMs: async () => (await getRequiredRadarProfile()).collectionIntervalMs, collect: collectScheduledSources, + getNextCycleAt: () => createDailyPublicationSchedule(() => new Date()).getNextPublicationAt(), onError: (error) => console.error("Task Worker 任务失败:", error), publish: publishDailyBriefWhenDue, - timers: { clearInterval, clearTimeout, setInterval, setTimeout }, + timers: { clearTimeout, setTimeout }, }); if (process.env.RADAR_WORKER_ONCE === "true") { if (await schedule.runOnce() === "failed") process.exitCode = 1; diff --git a/test/brief-publication.test.ts b/test/brief-publication.test.ts index c8bae1f..9b1a8f2 100644 --- a/test/brief-publication.test.ts +++ b/test/brief-publication.test.ts @@ -58,14 +58,15 @@ class InMemoryPublicationArchive implements PublicationArchive { this.delayedCandidates.push(input); } - async hasPublishedBrief(publicationDay: string) { - return this.publishedDays.includes(publicationDay); + async hasPublishedBrief(publicationDay: string, publicationSlot = "morning") { + return this.publishedDays.includes(`${publicationDay}:${publicationSlot}`); } async publishBrief(input: Parameters[0]) { - if (this.publishedDays.includes(input.publicationDay)) return "already-published" as const; + const key = `${input.publicationDay}:${input.publicationSlot}`; + if (this.publishedDays.includes(key)) return "already-published" as const; this.published = input; - this.publishedDays.push(input.publicationDay); + this.publishedDays.push(key); return "published" as const; } @@ -103,6 +104,7 @@ test("固定 Runtime 通过质量门后发布含 Section Citation 与 Provenance id: "brief-1", publishedAt: "2026-08-12T01:00:00.000Z", publicationDay: "2026-08-12", + publicationSlot: "morning", provenance: { configurationVersion: "profile@v1", modelRuntimeId: "compatible:fixture", @@ -138,6 +140,30 @@ test("固定 Runtime 通过质量门后发布含 Section Citation 与 Provenance ]); }); +test("同一 CST 日期可以分别发布早间和下午 Brief Snapshot", async () => { + const archive = new InMemoryPublicationArchive([candidate]); + const morning = await createFixturePublisher({ + archive, + clock: () => new Date("2026-08-12T01:00:00.000Z"), + createBriefId: () => "brief-morning", + isCitationAccessible: async () => true, + publicationSlot: "morning", + runtime: runtimeFor(validAssessment), + }).publishDailyBrief(); + const afternoon = await createFixturePublisher({ + archive, + clock: () => new Date("2026-08-12T09:00:00.000Z"), + createBriefId: () => "brief-afternoon", + isCitationAccessible: async () => true, + publicationSlot: "afternoon", + runtime: runtimeFor(validAssessment), + }).publishDailyBrief(); + + assert.equal(morning.status, "published"); + assert.equal(afternoon.status, "published"); + assert.deepEqual(archive.publishedDays, ["2026-08-12:morning", "2026-08-12:afternoon"]); +}); + test("日报只消费持久化完成的评估,不会再次调用模型", async () => { const archive = new InMemoryPublicationArchive([]); const evidenceExcerpt = "Codex 是一个帮助 Builder 在本地代码任务中完成实现、验证与审查的开源工具。"; @@ -266,7 +292,7 @@ test("日报发布在同一 CST 日期幂等,跨日创建新的不可变 Snaps assert.deepEqual(await publish("2026-08-12T01:00:00.000Z", "brief-0812"), { briefId: "brief-0812", signalCount: 1, status: "published" }); assert.deepEqual(await publish("2026-08-12T04:00:00.000Z", "brief-0812-again"), { status: "already-published" }); assert.deepEqual(await publish("2026-08-13T01:00:00.000Z", "brief-0813"), { briefId: "brief-0813", signalCount: 1, status: "published" }); - assert.deepEqual(archive.publishedDays, ["2026-08-12", "2026-08-13"]); + assert.deepEqual(archive.publishedDays, ["2026-08-12:morning", "2026-08-13:morning"]); }); test("校验或运行时失败会留下可诊断的阶段记录,且不发布 Snapshot", async () => { diff --git a/test/browser/mobile-reading.spec.ts b/test/browser/mobile-reading.spec.ts index a21f5e8..fc947d2 100644 --- a/test/browser/mobile-reading.spec.ts +++ b/test/browser/mobile-reading.spec.ts @@ -212,7 +212,8 @@ test("Instance Administrator 可在紧凑屏幕加载并编辑 Profile,未授 await page.getByRole("button", { name: "加载配置" }).click(); await expect(page.getByText("来源连接器")).toBeVisible(); - await page.getByRole("spinbutton", { name: "采集间隔(分钟)" }).fill("120"); + await expect(page.getByRole("textbox", { name: "自动采集时间" })).toHaveValue("每天 09:00、17:00(中国标准时间)"); + await page.getByRole("spinbutton", { name: "每轮评估上限" }).fill("6"); await page.getByRole("button", { name: "校验并保存新版本" }).click(); await expect(page.getByText("Compatible API 凭据未由部署环境配置。")).toBeVisible(); }); diff --git a/test/daily-publication-schedule.test.ts b/test/daily-publication-schedule.test.ts index 7b18701..1934563 100644 --- a/test/daily-publication-schedule.test.ts +++ b/test/daily-publication-schedule.test.ts @@ -2,14 +2,23 @@ import assert from "node:assert/strict"; import test from "node:test"; import { createDailyPublicationSchedule } from "../src/lib/radar/daily-publication-schedule.ts"; -test("固定时钟以 CST 09:00 作为日报发布边界,并给出下一次发布时刻", () => { +test("固定时钟以 CST 09:00 和 17:00 作为两个日报发布边界", () => { const beforeNine = createDailyPublicationSchedule(() => new Date("2026-08-12T00:59:59.000Z")); assert.equal(beforeNine.getDuePublicationDay(), null); assert.equal(beforeNine.getNextPublicationAt().toISOString(), "2026-08-12T01:00:00.000Z"); const atNine = createDailyPublicationSchedule(() => new Date("2026-08-12T01:00:00.000Z")); + assert.deepEqual(atNine.getDuePublication(), { day: "2026-08-12", slot: "morning" }); assert.equal(atNine.getDuePublicationDay(), "2026-08-12"); - assert.equal(atNine.getNextPublicationAt().toISOString(), "2026-08-13T01:00:00.000Z"); + assert.equal(atNine.getNextPublicationAt().toISOString(), "2026-08-12T09:00:00.000Z"); + + const atFive = createDailyPublicationSchedule(() => new Date("2026-08-12T09:00:00.000Z")); + assert.deepEqual(atFive.getDuePublications(), [ + { day: "2026-08-12", slot: "morning" }, + { day: "2026-08-12", slot: "afternoon" }, + ]); + assert.deepEqual(atFive.getDuePublication(), { day: "2026-08-12", slot: "afternoon" }); + assert.equal(atFive.getNextPublicationAt().toISOString(), "2026-08-13T01:00:00.000Z"); }); test("重启后的当日午后仍指向同一个 CST 发布日", () => { diff --git a/test/e2e/brief-publication.test.ts b/test/e2e/brief-publication.test.ts index 6bae30b..0de875f 100644 --- a/test/e2e/brief-publication.test.ts +++ b/test/e2e/brief-publication.test.ts @@ -293,6 +293,41 @@ test("固定 Runtime 经真实 PostgreSQL 按日发布后,API 读取 Snapshot assert.deepEqual(brief.signals[0]?.evidence, [{ label: "openai/codex", source: "GitHub Trending", url: "https://github.com/openai/codex" }]); }); +test("真实 PostgreSQL 允许同一 CST 日期分别发布早间和下午 Snapshot", { concurrency: false }, async () => { + const morning = await createBriefPublisher({ + archive: postgresBriefPublicationArchive, + clock: () => new Date("2026-08-12T01:00:00.000Z"), + configurationVersion: "profile@v1", + createBriefId: () => "brief-e2e-morning", + isCitationAccessible: async () => true, + pipelineVersion: "assessment-pipeline@v1", + publicationSlot: "morning", + runtime: fixedRuntime(validAssessment), + }).publishDailyBrief(); + assert.deepEqual(morning, { briefId: "brief-e2e-morning", signalCount: 1, status: "published" }); + + await seedSecondCandidate(); + const afternoon = await createBriefPublisher({ + archive: postgresBriefPublicationArchive, + clock: () => new Date("2026-08-12T09:00:00.000Z"), + configurationVersion: "profile@v1", + createBriefId: () => "brief-e2e-afternoon", + isCitationAccessible: async () => true, + pipelineVersion: "assessment-pipeline@v1", + publicationSlot: "afternoon", + runtime: candidateAwareRuntime("compatible:fixed-e2e-afternoon"), + }).publishDailyBrief(); + assert.deepEqual(afternoon, { briefId: "brief-e2e-afternoon", signalCount: 1, status: "published" }); + + const snapshots = await getDatabasePool().query<{ publication_day: string; publication_slot: string }>( + "SELECT publication_day::text AS publication_day, publication_slot FROM brief_snapshots ORDER BY published_at", + ); + assert.deepEqual(snapshots.rows, [ + { publication_day: "2026-08-12", publication_slot: "morning" }, + { publication_day: "2026-08-12", publication_slot: "afternoon" }, + ]); +}); + test("Daily Brief 冻结发布时的来源覆盖度,不读取后续实时 Connector Health", { concurrency: false }, async () => { const database = getDatabasePool(); await database.query( diff --git a/test/e2e/task-queue.test.ts b/test/e2e/task-queue.test.ts index 0ee7bf1..f2f3d3e 100644 --- a/test/e2e/task-queue.test.ts +++ b/test/e2e/task-queue.test.ts @@ -144,6 +144,13 @@ test("模型任务三次失败后持久化为评估延迟,并保留错误和 assert.deepEqual(taskState.rows, [{ attempt_count: 3, last_error: "Ollama Runtime 请求失败:HTTP 503", status: "delayed" }]); assert.equal(statistics.pendingByState["评估延迟"], 1); assert.equal(statistics.retryCount, 1); + + const requeued = await postgresCandidateTaskArchive.requeueDelayedAssessments?.({ configurationVersion: "profile@v2", runtimeId: "ollama:qwen3-local:8b" }); + assert.equal(requeued, 1); + const recoveredState = await getDatabasePool().query<{ assessment_delay_detail: string | null; evaluation_status: string }>("SELECT evaluation_status, assessment_delay_detail FROM radar_candidates WHERE id = $1", [input.canonicalIdentifier]); + assert.deepEqual(recoveredState.rows, [{ assessment_delay_detail: null, evaluation_status: "queued" }]); + const retryTask = await getDatabasePool().query<{ configuration_version: string; runtime_id: string; status: string; task_kind: string }>("SELECT task_kind, status, configuration_version, runtime_id FROM candidate_tasks WHERE status = 'queued'"); + assert.deepEqual(retryTask.rows, [{ configuration_version: "profile@v2", runtime_id: "ollama:qwen3-local:8b", status: "queued", task_kind: "enrichment" }]); }); test("评估任务持久关联其实际使用的 Primary Evidence", { concurrency: false }, async () => { diff --git a/test/profile-collection.test.ts b/test/profile-collection.test.ts index 5ba8b52..b8b8989 100644 --- a/test/profile-collection.test.ts +++ b/test/profile-collection.test.ts @@ -4,7 +4,6 @@ import { createProfileCandidateFilter, createProfileSourceConnectors } from "../ import type { RadarProfile } from "../src/lib/radar/radar-profile.ts"; const profile: RadarProfile = { - collectionIntervalMs: 7_200_000, enabledConnectorIds: ["github-trending", "show-hn"], excludeTerms: ["ignore"], id: "profile@v2", diff --git a/test/profile-route.test.ts b/test/profile-route.test.ts index 928dc0d..2712871 100644 --- a/test/profile-route.test.ts +++ b/test/profile-route.test.ts @@ -4,7 +4,6 @@ import { createProfileCollectionHandler, createProfileGetHandler, createProfileM import type { RadarProfile } from "../src/lib/radar/radar-profile.ts"; const profile: RadarProfile = { - collectionIntervalMs: 7_200_000, enabledConnectorIds: ["github-trending"], excludeTerms: [], id: "profile@v1", @@ -61,11 +60,11 @@ test("首次配置时,Profile API 返回可编辑草稿而不启用未校验 test("Profile API 将无效的初始环境草稿转为明确错误", async () => { const dependency = dependencies(); dependency.getActive = async () => null; - dependency.getDraft = () => { throw new Error("RADAR_COLLECTION_INTERVAL_MS 必须是有效整数。"); }; + dependency.getDraft = () => { throw new Error("Radar Profile 配置无效。"); }; const get = createProfileGetHandler(dependency); const response = await get(request("GET")); assert.equal(response.status, 400); - assert.deepEqual(await response.json(), { error: "RADAR_COLLECTION_INTERVAL_MS 必须是有效整数。" }); + assert.deepEqual(await response.json(), { error: "Radar Profile 配置无效。" }); }); test("Profile 回滚会验证目标版本的运行时后才激活", async () => { @@ -98,6 +97,16 @@ test("管理员可显式复评当前已评估待发布集合", async () => { assert.equal(requeued, 1); }); +test("管理员可重新排队评估延迟 Candidate", async () => { + const retry = createProfileReassessHandler({ + environment: { RADAR_ADMIN_TOKEN: "admin-token" }, + requeue: async () => 4, + }); + + assert.equal((await retry(request("POST", undefined, ""))).status, 401); + assert.deepEqual(await (await retry(request("POST"))).json(), { requeuedCount: 4 }); +}); + test("管理员可立即采集,但不会触发 Daily Brief 发布", async () => { let collectionRuns = 0; const collect = createProfileCollectionHandler({ diff --git a/test/radar-brief.test.ts b/test/radar-brief.test.ts index 681ae10..e7a96cc 100644 --- a/test/radar-brief.test.ts +++ b/test/radar-brief.test.ts @@ -105,9 +105,11 @@ test("Brief API 将 Assessment Delay 与已有已发布日报同时透明呈现" candidateCount: 2, detail: "Compatible Runtime 请求失败:HTTP 503(已重试 3 次)", }, - availability: "assessment-delayed", + availability: "published", signals: [publishedSignal], }); + + assert.equal(getAssessmentBanner({ ...payload, hasPublishedSignals: true, visibleSignalCount: 1 }), "另有 2 个 Candidate 评估延迟,当前已发布日报正常展示。原因:Compatible Runtime 请求失败:HTTP 503(已重试 3 次)"); }); test("未配置数据库时,公开 Brief 不再回退到 Fixture", async () => { diff --git a/test/radar-profile.test.ts b/test/radar-profile.test.ts index bfcbabc..625eb20 100644 --- a/test/radar-profile.test.ts +++ b/test/radar-profile.test.ts @@ -3,7 +3,6 @@ import test from "node:test"; import { createInitialRadarProfileConfig, parseRadarProfileConfig } from "../src/lib/radar/radar-profile.ts"; const profile = { - collectionIntervalMs: 7_200_000, enabledConnectorIds: ["github-trending", "show-hn"], excludeTerms: ["irrelevant"], includeTerms: ["agent"], @@ -31,7 +30,6 @@ test("Radar Profile 只接受受限的非敏感运行时配置与三个内置来 test("初始 Profile 从既有部署环境迁移,凭据不会进入配置", () => { const initial = createInitialRadarProfileConfig({ - RADAR_COLLECTION_INTERVAL_MS: "3600000", RADAR_INCLUDE_TERMS: "Agent, 推理", RADAR_MODEL_RUNTIME: "ollama", RADAR_OLLAMA_BASE_URL: "http://127.0.0.1:11434", diff --git a/test/task-queue.test.ts b/test/task-queue.test.ts index 252cd1d..de0944e 100644 --- a/test/task-queue.test.ts +++ b/test/task-queue.test.ts @@ -223,3 +223,46 @@ test("模型报告证据不足时完成任务而不进入运行时失败重试", assert.deepEqual(completed, ["assessment-insufficient"]); assert.deepEqual(failures, []); }); + +test("运行时不匹配的遗留评估任务不会阻塞同一周期的补证任务", async () => { + const staleAssessmentTask: ClaimedCandidateTask = { + ...enrichmentTask, + id: "stale-assessment", + kind: "assessment", + runtimeId: "legacy", + candidate: { + ...enrichmentTask.candidate, + primaryEvidence: [{ canonicalIdentifier: "primary:github", contentFingerprint: "a".repeat(64), excerpts: ["A local coding agent."], fetchedAt: enrichmentTask.claimedAt, sourceKind: "github-repository-description", sourceName: "GitHub", sourceTitle: "codex", sourceUrl: enrichmentTask.candidate.url }], + }, + }; + const enrichmentTaskToRun = { ...enrichmentTask, id: "enrichment-after-stale-task" }; + const tasks = [staleAssessmentTask, enrichmentTaskToRun]; + const released: string[] = []; + const completed: string[] = []; + const archive: CandidateTaskArchive = { + claimNext: async ({ excludeTaskIds }) => { + const index = tasks.findIndex((task) => !excludeTaskIds?.includes(task.id)); + return index < 0 ? null : tasks[index]; + }, + completeAssessment: async () => undefined, + completeEnrichment: async ({ task }) => { completed.push(task.id); tasks.splice(tasks.findIndex((item) => item.id === task.id), 1); }, + enqueueEnrichment: async () => undefined, + fail: async () => "retryable", + getStatistics: async () => ({ averageDurationMs: 0, completedThisCycle: 0, estimatedDrainCount: 0, estimatedDrainMs: 0, pendingByState: { "待补证": 0, "补证中": 0, "评估中": 0, "评估失败待重试": 0, "评估延迟": 0, "证据不足未入选": 0, "已评估未入选": 0, "已评估待发布": 0 }, retryCount: 0 }), + release: async ({ task }) => { released.push(task.id); }, + requeueReadyAssessments: async () => 0, + }; + const worker = createCandidateTaskWorker({ + archive, + clock: () => new Date("2026-08-13T01:00:00.000Z"), + enrich: async () => ({ candidateCanonicalIdentifier: enrichmentTask.candidate.canonicalIdentifier, digests: [], status: "insufficient-evidence" as const }), + getRuntime: async () => ({ assess: async () => ({ assessmentOutcome: "insufficient-evidence", assessmentReason: "不应评估遗留任务。" } as never), id: "ollama:qwen3" }), + maxTasks: 2, + timeBudgetMs: 10_000, + workerId: "worker-1", + }); + + assert.equal(await worker.runCycle(), 2); + assert.deepEqual(released, ["stale-assessment"]); + assert.deepEqual(completed, ["enrichment-after-stale-task"]); +}); diff --git a/test/task-worker-schedule.test.ts b/test/task-worker-schedule.test.ts index d39df4c..87b7172 100644 --- a/test/task-worker-schedule.test.ts +++ b/test/task-worker-schedule.test.ts @@ -6,21 +6,13 @@ type ScheduledTask = { callback: () => void; delay: number; id: number }; function createFakeTimers() { let nextId = 1; - const intervals: ScheduledTask[] = []; const timeouts: ScheduledTask[] = []; const cleared: number[] = []; return { cleared, - intervals, timeouts, timers: { - clearInterval: (id: number) => cleared.push(id), clearTimeout: (id: number) => cleared.push(id), - setInterval: (callback: () => void, delay: number) => { - const task = { callback, delay, id: nextId++ }; - intervals.push(task); - return task.id; - }, setTimeout: (callback: () => void, delay: number) => { const task = { callback, delay, id: nextId++ }; timeouts.push(task); @@ -30,33 +22,33 @@ function createFakeTimers() { }; } -test("Worker 启动后每两小时采集,并在 CST 09:00 发布后继续安排下一日", async () => { +test("Worker 在 CST 09:00 和 17:00 各采集并发布一次", async () => { const fake = createFakeTimers(); const events: string[] = []; let now = new Date("2026-08-12T00:59:59.000Z"); const schedule = createTaskWorkerSchedule({ clock: () => now, - collectionIntervalMs: 2 * 60 * 60 * 1000, collect: async () => { events.push("collect"); return "succeeded" as const; }, + getNextCycleAt: () => new Date("2026-08-12T01:00:00.000Z"), publish: async () => { events.push("publish"); }, timers: fake.timers, }); await schedule.start(); assert.deepEqual(events, ["collect"]); - assert.deepEqual(fake.timeouts.map(({ delay }) => delay), [7_200_000, 1_000]); + assert.deepEqual(fake.timeouts.map(({ delay }) => delay), [1_000]); now = new Date("2026-08-12T01:00:00.000Z"); - fake.timeouts[1]?.callback(); + fake.timeouts[0]?.callback(); await new Promise((resolve) => setImmediate(resolve)); - assert.deepEqual(events, ["collect", "publish"]); - assert.deepEqual(fake.timeouts.map(({ delay }) => delay), [7_200_000, 1_000, 86_400_000]); + assert.deepEqual(events, ["collect", "collect", "publish"]); - fake.timeouts[0]?.callback(); + now = new Date("2026-08-12T09:00:00.000Z"); + fake.timeouts[1]?.callback(); await new Promise((resolve) => setImmediate(resolve)); - assert.deepEqual(events, ["collect", "publish", "collect"]); + assert.deepEqual(events, ["collect", "collect", "publish", "collect", "publish"]); schedule.stop(); - assert.deepEqual(fake.cleared, [4, 3]); + assert.deepEqual(fake.cleared, [3]); }); test("Worker 重启于 CST 09:00 后会补发当天日报", async () => { @@ -64,8 +56,8 @@ test("Worker 重启于 CST 09:00 后会补发当天日报", async () => { const events: string[] = []; const schedule = createTaskWorkerSchedule({ clock: () => new Date("2026-08-12T06:00:00.000Z"), - collectionIntervalMs: 2 * 60 * 60 * 1000, collect: async () => { events.push("collect"); return "succeeded" as const; }, + getNextCycleAt: () => new Date("2026-08-12T09:00:00.000Z"), publish: async () => { events.push("publish"); }, timers: fake.timers, }); @@ -73,7 +65,7 @@ test("Worker 重启于 CST 09:00 后会补发当天日报", async () => { await schedule.start(); assert.deepEqual(events, ["collect", "publish"]); - assert.deepEqual(fake.timeouts.map(({ delay }) => delay), [7_200_000, 68_400_000]); + assert.deepEqual(fake.timeouts.map(({ delay }) => delay), [10_800_000]); }); test("定时任务遇到采集或发布异常会记录错误,并继续安排下一次执行", async () => { @@ -82,8 +74,8 @@ test("定时任务遇到采集或发布异常会记录错误,并继续安排 let now = new Date("2026-08-12T00:59:59.000Z"); const schedule = createTaskWorkerSchedule({ clock: () => now, - collectionIntervalMs: 2 * 60 * 60 * 1000, collect: async () => { throw new Error("采集网络异常"); }, + getNextCycleAt: () => new Date("2026-08-12T01:00:00.000Z"), onError: (error) => { errors.push(error); }, publish: async () => { throw new Error("发布数据库异常"); }, timers: fake.timers, @@ -93,23 +85,19 @@ test("定时任务遇到采集或发布异常会记录错误,并继续安排 assert.equal((errors[0] as Error)?.message, "采集网络异常"); now = new Date("2026-08-12T01:00:00.000Z"); - fake.timeouts[1]?.callback(); - await new Promise((resolve) => setImmediate(resolve)); - assert.deepEqual(fake.timeouts.map(({ delay }) => delay), [7_200_000, 1_000, 86_400_000]); - fake.timeouts[0]?.callback(); await new Promise((resolve) => setImmediate(resolve)); - assert.deepEqual(errors.map((error) => (error as Error).message), ["采集网络异常", "发布数据库异常", "采集网络异常"]); + assert.deepEqual(fake.timeouts.map(({ delay }) => delay), [1_000, 0]); + assert.deepEqual(errors.map((error) => (error as Error).message), ["采集网络异常", "采集网络异常"]); }); -test("Worker 在下一 Collection Cycle 读取新的采集间隔", async () => { +test("Worker 在下一时段读取新的采集时间", async () => { const fake = createFakeTimers(); - const intervals = [120_000, 240_000]; + const nextCycles = [new Date("2026-08-12T00:02:00.000Z"), new Date("2026-08-12T00:04:00.000Z")]; const schedule = createTaskWorkerSchedule({ clock: () => new Date("2026-08-12T00:00:00.000Z"), - collectionIntervalMs: 60_000, collect: async () => "succeeded" as const, - getCollectionIntervalMs: async () => intervals.shift() ?? 240_000, + getNextCycleAt: () => nextCycles.shift() ?? new Date("2026-08-12T01:04:00.000Z"), publish: async () => undefined, timers: fake.timers, }); @@ -118,25 +106,27 @@ test("Worker 在下一 Collection Cycle 读取新的采集间隔", async () => { assert.equal(fake.timeouts[0]?.delay, 120_000); fake.timeouts[0]?.callback(); await new Promise((resolve) => setImmediate(resolve)); - assert.equal(fake.timeouts[2]?.delay, 240_000); + assert.equal(fake.timeouts[1]?.delay, 240_000); }); test("Worker 停止后不会由已完成的发布回调重新登记定时器", async () => { const fake = createFakeTimers(); + let now = new Date("2026-08-12T00:59:59.000Z"); let resolvePublish: (() => void) | undefined; const schedule = createTaskWorkerSchedule({ - clock: () => new Date("2026-08-12T00:59:59.000Z"), - collectionIntervalMs: 60_000, + clock: () => now, collect: async () => "succeeded" as const, + getNextCycleAt: () => new Date("2026-08-12T01:00:00.000Z"), publish: async () => new Promise((resolve) => { resolvePublish = resolve; }), timers: fake.timers, }); await schedule.start(); - fake.timeouts[1]?.callback(); + now = new Date("2026-08-12T01:00:00.000Z"); + fake.timeouts[0]?.callback(); await new Promise((resolve) => setImmediate(resolve)); schedule.stop(); resolvePublish?.(); await new Promise((resolve) => setImmediate(resolve)); - assert.equal(fake.timeouts.length, 2); + assert.equal(fake.timeouts.length, 1); });