diff --git a/.gitignore b/.gitignore
index e098e89e..f4ace5a5 100644
--- a/.gitignore
+++ b/.gitignore
@@ -73,3 +73,6 @@ htmlcov/
.brooks-lint-history.json
.vidt/
.node-version
+
+# 独立 worktree 工作区(issue 切片隔离)
+.worktrees/
diff --git a/backend/.env.example b/backend/.env.example
index 91ce7cfc..5909e26c 100644
--- a/backend/.env.example
+++ b/backend/.env.example
@@ -70,6 +70,25 @@ AUTH_REFRESH_ATTEMPTS_PER_MINUTE=20
LLM_REQUESTS_PER_MINUTE=30
# 单次模型调用硬截止时间;模型配置中的 extra_params.litellm_params.timeout 只能缩短。
LLM_COMPLETION_TIMEOUT_SECONDS=45
+# 预算熔断:按实际 DONE 调用数的三窗保险丝(分钟/小时/日),超限拒绝新付费调用,
+# 分析链路自动走本地降级(卡片标「本地速览·待 AI 复核」),预算余量 >30% 时由
+# requeue job 自动补分析;窗口滚过后恢复。0 = 关闭对应窗口(代码默认全关,opt-in)。
+LLM_BUDGET_CALLS_PER_MINUTE=60
+LLM_BUDGET_CALLS_PER_HOUR=2000
+LLM_BUDGET_CALLS_PER_DAY=20000
+#
+# 数值校准依据(2026-09-30 圆桌,四席对照;实测负载:日峰 5865 / 时峰 508 / 分峰 24):
+# 现阶段维持 60/2000/20000——圆桌裁决:requeue 补分析刚落地(#90)、
+# 尚未稳定运行观察,误触发代价(本地降级)未完全可逆前天平向宽松;
+# 灾难场景仍有 3.4x+ 余量。
+# SRE 席建议收紧 70/1500/10000-12000:限流 30/min 下时窗 2000 数学上不可达(持续上限 1800/h);
+# 成本席建议收紧 60/1200/10000:全旗舰路由最坏 ¥240-3600/日,次数上限锁不住金额上限;
+# 产品席偏松 60/3000/20000:误触发=降级,代价不可逆时宁可松;
+# 数据席维持 60/2000/20000:现值对正常流量零误伤;relation_discovery 9 天涨 3.6 倍,日窗最怕增长而非异常。
+# 收紧前提:requeue 补分析(已随 #90 落地)稳定运行后,可收紧至 60/1200/10000 并加 10k 告警线;
+# 切换付费 API 后建议二期增加金额三窗(当前 total_cost 记账未实现,暂无数据可依)。
+# 预算豁免场景(逗号分隔):低频核心承诺场景不做预算检查,防日报在日窗耗尽时硬失败。
+LLM_BUDGET_EXEMPT_SCENES=daily_report,weekly_digest,monthly_digest
# OAuth (Google / GitHub 登录)
# 留空则该 provider 不启用。前端登录页会根据已启用 provider 渲染对应按钮。
diff --git a/backend/app/core/config.py b/backend/app/core/config.py
index 9738c6f2..84552c2c 100644
--- a/backend/app/core/config.py
+++ b/backend/app/core/config.py
@@ -176,6 +176,16 @@ class Settings(BaseSettings):
LLM_WORKER_CONCURRENCY: int = 4
# 单次运行时调用的硬上限;模型 extra_params.timeout 只能把它调小,不能放大。
LLM_COMPLETION_TIMEOUT_SECONDS: float = 45.0
+ # 预算熔断(budget_guard):按 llm_call_logs 实际 DONE 调用数的三窗保险丝,
+ # 超限直接拒绝新调用(分析链路走本地降级,窗口过后自动恢复)。
+ # 默认 0 = 全部关闭(opt-in,不改变现有部署行为);0 同样适用于单窗关闭。
+ # 开启建议:分钟窗不低于 LLM_REQUESTS_PER_MINUTE,日窗按可接受的烧钱上限设。
+ LLM_BUDGET_CALLS_PER_MINUTE: int = 0
+ LLM_BUDGET_CALLS_PER_HOUR: int = 0
+ LLM_BUDGET_CALLS_PER_DAY: int = 0
+ # 预算豁免场景(逗号分隔):这些 scene 的调用不做预算检查。日报/周报/月报
+ # 低频且是核心承诺(预算耗尽会让 daily_report 直接走 ERROR),默认豁免。
+ LLM_BUDGET_EXEMPT_SCENES: str = "daily_report,weekly_digest,monthly_digest"
ANALYSIS_WORKER_CONCURRENCY: int = 3
ANALYSIS_MAX_ATTEMPTS: int = 5
ANALYSIS_RETRY_BASE_DELAY_SECONDS: int = 60
diff --git a/backend/app/repositories/content_repo.py b/backend/app/repositories/content_repo.py
index 3b7abbea..f93031c9 100644
--- a/backend/app/repositories/content_repo.py
+++ b/backend/app/repositories/content_repo.py
@@ -783,6 +783,77 @@ async def list_favorites(
# ── Today picks candidates (SQLite fallback) ─────────────────
+ async def reset_local_fallback_for_reanalysis(
+ self,
+ *,
+ limit: int = 50,
+ cooldown_minutes: int = 60,
+ max_age_days: int = 7,
+ ) -> int:
+ """Requeue local_fallback-analyzed content for real LLM analysis (#90).
+
+ 降级触发原因(预算窗口/熔断)消除后由 requeue job 调用,把永久假分
+ 变回真实分析。只回收「最新一条分析是 local_fallback」的 ANALYZED 内容:
+
+ - 最旧优先(``limit`` 限速,避免一次灌满分析队列);
+ - ``cooldown_minutes``:降级落库后至少等待,防止预算临界值附近
+ requeue→再降级→再 requeue 的反复横跳;
+ - ``max_age_days``:只回收近 N 天(旧内容重分析价值低);
+ - ``skip_analysis`` 为 True 的内容不回收(显式跳过 LLM 的语义优先);
+ - 已有 ≥2 条 local_fallback 分析的内容不回收:attempts=0 重置使
+ ANALYSIS_MAX_ATTEMPTS 失效,内容过滤类确定性降级(不消耗预算
+ 计数、余量门控拦不住)反复 requeue 会浪费调用(#90 复核结论)。
+
+ 返回实际重置条数。SELECT+UPDATE 两步:UPDATE 复核 status=ANALYZED
+ 防两步间隙状态漂移。
+ """
+ from datetime import timedelta
+
+ from sqlalchemy import update
+
+ from app.models.analysis import AiAnalysis
+
+ now = naive_utc_now()
+ latest_analysis_id = self._latest_analysis_id_subquery(AiAnalysis)
+ candidate_ids_stmt = (
+ select(self.model.id)
+ .join(AiAnalysis, AiAnalysis.id == latest_analysis_id)
+ .where(self.model.status == ContentStatus.ANALYZED)
+ .where(self.model.skip_analysis.is_(False))
+ .where(self.model.updated_at <= now - timedelta(minutes=cooldown_minutes))
+ .where(self.model.created_at >= now - timedelta(days=max_age_days))
+ .where(AiAnalysis.summary_source == "local_fallback")
+ .where(
+ self.model.id.notin_(
+ select(AiAnalysis.content_id)
+ .where(AiAnalysis.summary_source == "local_fallback")
+ .group_by(AiAnalysis.content_id)
+ .having(func.count() >= 2)
+ )
+ )
+ .order_by(self.model.created_at.asc())
+ .limit(limit)
+ )
+ rows = await self.db.execute(candidate_ids_stmt)
+ ids = [int(row[0]) for row in rows.all()]
+ if not ids:
+ return 0
+
+ result = await self.db.execute(
+ update(self.model)
+ .where(self.model.id.in_(ids))
+ .where(self.model.status == ContentStatus.ANALYZED)
+ .values(
+ status=ContentStatus.PENDING,
+ updated_at=now,
+ analysis_attempts=0,
+ analysis_next_retry_at=None,
+ analysis_claim_token=None,
+ analysis_lease_expires_at=None,
+ )
+ )
+ return int(result.rowcount or 0)
+
async def list_for_today_picks(
self,
hours: int = 48,
@@ -808,7 +879,9 @@ async def list_for_today_picks(
selectinload(self.model.source),
)
.join(AiAnalysis, AiAnalysis.id == latest_analysis_id)
- .where(self.model.crawled_at >= cutoff)
+ # 窗口按原文发布时间优先(published_at 为空回退抓取时间),旧文晚抓
+ # 不进「今天」;与 DuckDB query_today_picks 的 COALESCE 口径对齐。
+ .where(func.coalesce(self.model.published_at, self.model.crawled_at) >= cutoff)
.where(AiAnalysis.risk_score <= risk_threshold)
)
if category:
@@ -847,8 +920,10 @@ async def list_for_report_window(
selectinload(self.model.source),
)
.join(AiAnalysis, AiAnalysis.id == latest_analysis_id)
- .where(self.model.crawled_at >= window_start)
- .where(self.model.crawled_at <= window_end)
+ # 报告窗口同样按原文发布时间归档:晚抓到的旧文归入其发布期,
+ # 不出现在之后的日报里(口径与 list_for_today_picks 一致)。
+ .where(func.coalesce(self.model.published_at, self.model.crawled_at) >= window_start)
+ .where(func.coalesce(self.model.published_at, self.model.crawled_at) <= window_end)
.where(AiAnalysis.risk_score <= risk_threshold)
.where(AiAnalysis.curation_score.isnot(None))
.where(~self._accepted_event_member_exists())
diff --git a/backend/app/scheduler.py b/backend/app/scheduler.py
index 8d5c1d3f..2ec70ad3 100644
--- a/backend/app/scheduler.py
+++ b/backend/app/scheduler.py
@@ -890,6 +890,20 @@ async def _normalize_content_events() -> dict:
# ── Lifecycle helpers ─────────────────────────────────────────────────
+async def _requeue_local_fallback() -> None:
+ """#90:预算/熔断恢复后补分析降级内容(余量门控见 analysis_requeue)。"""
+ from app.services.analysis_requeue import requeue_local_fallback_content
+
+ try:
+ async with async_session() as db:
+ count = await requeue_local_fallback_content(db)
+ await db.commit()
+ if count:
+ logger.info("Scheduler: requeued %d local_fallback item(s) for re-analysis", count)
+ except Exception:
+ logger.exception("Scheduler: local_fallback requeue failed")
+
+
def start_scheduler() -> None:
"""Register all scheduled jobs and start the scheduler."""
if scheduler.running:
@@ -1113,6 +1127,16 @@ def start_scheduler() -> None:
replace_existing=True,
)
+ # 降级内容补分析 (#90):每 15 分钟在预算余量充足时,把 local_fallback
+ # 降级内容重置回 PENDING 让分析队列重新拿真实 LLM 分析
+ scheduler.add_job(
+ _requeue_local_fallback,
+ trigger=IntervalTrigger(minutes=15),
+ id="requeue_local_fallback",
+ name="降级内容补分析(预算余量门控)",
+ replace_existing=True,
+ )
+
# 趋势快照:每日04:30计算话题热度+关键词频率
# (错开 04:00 知乎抓取 + 04:15 聚类,避免资源竞争)
scheduler.add_job(
diff --git a/backend/app/services/_duckdb_picks_mixin.py b/backend/app/services/_duckdb_picks_mixin.py
index 23e0e070..4a5744f7 100644
--- a/backend/app/services/_duckdb_picks_mixin.py
+++ b/backend/app/services/_duckdb_picks_mixin.py
@@ -105,7 +105,7 @@ def query_today_picks(
LEFT JOIN oltp_db.sources s ON s.id = c.source_id
LEFT JOIN feedback_scores f ON f.content_id = c.id
LEFT JOIN ignored_content ignored ON ignored.content_id = c.id
- WHERE c.crawled_at >= ?
+ WHERE COALESCE(c.published_at, c.crawled_at) >= ?
AND ignored.content_id IS NULL
AND a.risk_score <= {risk_threshold}
AND a.curation_score IS NOT NULL
diff --git a/backend/app/services/_duckdb_reports_mixin.py b/backend/app/services/_duckdb_reports_mixin.py
index 7af242fb..6b85d01a 100644
--- a/backend/app/services/_duckdb_reports_mixin.py
+++ b/backend/app/services/_duckdb_reports_mixin.py
@@ -29,7 +29,8 @@ def query_content_for_report(self, hours: int = 48) -> list[dict[str, Any]]:
FROM oltp_db.content_items c
LEFT JOIN latest_analysis a ON a.content_id = c.id
LEFT JOIN ignored_content ignored ON ignored.content_id = c.id
- WHERE c.crawled_at >= '{cutoff}'
+ -- 窗口按原文发布时间优先(口径同 list_for_report_window / query_today_picks)
+ WHERE COALESCE(c.published_at, c.crawled_at) >= '{cutoff}'
AND ignored.content_id IS NULL
AND a.curation_score IS NOT NULL
ORDER BY (COALESCE(a.creator_score, 0) + COALESCE(a.viral_score, 0)) DESC
@@ -79,14 +80,16 @@ def query_content_for_weekly(self, start_date: str, end_date: str) -> list[dict[
COALESCE(f.feedback_score, 0) AS feedback_score,
COALESCE(a.curation_score, 0)
+ LEAST({feedback_max}, GREATEST({feedback_min}, COALESCE(f.feedback_score, 0))) * {feedback_weight}
- AS adjusted_score
+ AS adjusted_score,
+ a.summary_source
FROM oltp_db.content_items c
LEFT JOIN latest_analysis a ON a.content_id = c.id
LEFT JOIN oltp_db.sources s ON s.id = c.source_id
LEFT JOIN feedback_scores f ON f.content_id = c.id
LEFT JOIN ignored_content ignored ON ignored.content_id = c.id
- WHERE CAST(c.crawled_at AS DATE) >= DATE '{start_date}'
- AND CAST(c.crawled_at AS DATE) <= DATE '{end_date}'
+ -- 周报/月报同样按原文发布时间归档,口径与日报窗口一致
+ WHERE CAST(COALESCE(c.published_at, c.crawled_at) AS DATE) >= DATE '{start_date}'
+ AND CAST(COALESCE(c.published_at, c.crawled_at) AS DATE) <= DATE '{end_date}'
AND ignored.content_id IS NULL
AND a.curation_score IS NOT NULL
ORDER BY adjusted_score DESC, COALESCE(a.creator_score, 0) DESC
@@ -118,6 +121,7 @@ def query_content_for_weekly(self, start_date: str, end_date: str) -> list[dict[
"source_weight_db": int(row[21]) if row[21] else 3,
"feedback_score": float(row[22]) if row[22] else 0,
"adjusted_score": round(float(row[23] or 0), 1),
+ "summary_source": row[24],
}
for row in results
]
diff --git a/backend/app/services/analysis.py b/backend/app/services/analysis.py
index 6530e4d6..4cc6a1eb 100644
--- a/backend/app/services/analysis.py
+++ b/backend/app/services/analysis.py
@@ -506,15 +506,19 @@ async def analyze_content(content: ContentItem, db: AsyncSession) -> AiAnalysis:
)
final_model = final_metadata.get("actual_model") or final_model
except Exception as llm_exc:
- # CircuitOpenError (breaker tripped) and BadRequestError (400,
- # e.g. GLM contentFilter code=1301) trigger local fallback.
- # Other LLM failures (timeout, network, RuntimeError) still
- # propagate up so the caller can record ERROR status + retry.
+ # CircuitOpenError (breaker tripped), LlmBudgetExceededError (budget
+ # fuse) and BadRequestError (400, e.g. GLM contentFilter code=1301)
+ # trigger local fallback. Other LLM failures (timeout, network,
+ # RuntimeError) still propagate up so the caller can record ERROR
+ # status + retry. Budget rejection must especially NOT consume
+ # analysis attempts: a day-window fuse can outlast the retry budget
+ # and would otherwise burn the whole backlog into permanent ERROR.
from litellm.exceptions import BadRequestError
+ from app.services.llm.budget_guard import LlmBudgetExceededError
from app.services.llm.circuit_breaker import CircuitOpenError
- if isinstance(llm_exc, CircuitOpenError | BadRequestError):
+ if isinstance(llm_exc, CircuitOpenError | BadRequestError | LlmBudgetExceededError):
logger.warning(
"LLM call failed for content id=%d (%s), using local fallback",
content.id,
diff --git a/backend/app/services/analysis_requeue.py b/backend/app/services/analysis_requeue.py
new file mode 100644
index 00000000..44499f36
--- /dev/null
+++ b/backend/app/services/analysis_requeue.py
@@ -0,0 +1,47 @@
+"""降级内容补分析(requeue)——#90 的核心修复。
+
+LLM 降级路径(熔断 / 内容过滤 / 预算三种触发源)曾把内容写成
+``ANALYZED`` 终态后永久结疤(假分 61-64 窄带、无任何消费者捡起)。
+本模块在触发原因消除后把 ``summary_source='local_fallback'`` 的内容
+重置回 PENDING,让既有分析队列重新拿真实 LLM 分析——把「永久疤」
+变成「一过性」。
+
+门控与限速(防预算临界值附近的反复横跳):
+
+- **预算余量门控**:最紧启用窗口余量 < 30% 时不发起(fail-closed,
+ 见 budget_guard.budget_headroom_ok);
+- **冷却**:降级落库后至少等 60 分钟才可回收;
+- **限量**:每轮最多 50 条,最旧优先。
+"""
+
+from __future__ import annotations
+
+import logging
+
+from sqlalchemy.ext.asyncio import AsyncSession
+
+from app.repositories.content_repo import ContentRepo
+from app.services.llm.budget_guard import budget_headroom_ok
+
+logger = logging.getLogger(__name__)
+
+
+async def requeue_local_fallback_content(
+ db: AsyncSession,
+ *,
+ limit: int = 50,
+ cooldown_minutes: int = 60,
+) -> int:
+ """预算余量充足时,把降级内容重置回待分析。返回重置条数。"""
+ if not await budget_headroom_ok(min_ratio=0.3):
+ logger.info("Analysis requeue: LLM budget headroom insufficient, deferring")
+ return 0
+
+ count = await ContentRepo(db).reset_local_fallback_for_reanalysis(
+ limit=limit,
+ cooldown_minutes=cooldown_minutes,
+ )
+ if count:
+ # 恢复是用户可感知的质量事件,按 warning 级别留痕(不是静默修复)
+ logger.warning("Analysis requeue: reset %d local_fallback content item(s) back to PENDING", count)
+ return count
diff --git a/backend/app/services/creation.py b/backend/app/services/creation.py
index 961eda4a..6457b5bc 100644
--- a/backend/app/services/creation.py
+++ b/backend/app/services/creation.py
@@ -26,6 +26,7 @@
def _creation_llm_error_message(exc: Exception) -> str:
"""Return an actionable, non-internal error message for creation callers."""
+ from app.services.llm.budget_guard import LlmBudgetExceededError
from app.services.llm.circuit_breaker import CircuitOpenError
from app.services.llm.provider import LlmCapacityUnavailableError
@@ -33,6 +34,8 @@ def _creation_llm_error_message(exc: Exception) -> str:
return "创作方案暂时排队等待可用模型渠道,请稍后重试。"
if isinstance(exc, CircuitOpenError):
return "创作方案服务暂时不可用,系统正在自动恢复,请稍后重试。"
+ if isinstance(exc, LlmBudgetExceededError):
+ return "模型调用预算已达今日上限,创作方案稍后自动恢复,请明日再试。"
return "创作方案生成失败,请稍后重试。"
diff --git a/backend/app/services/digest_context.py b/backend/app/services/digest_context.py
index b13d565c..7304c862 100644
--- a/backend/app/services/digest_context.py
+++ b/backend/app/services/digest_context.py
@@ -67,6 +67,7 @@ def value_or_default(value, default):
category=row.get("category"),
source_name=row.get("source_name"),
crawled_at=row.get("crawled_at"),
+ summary_source=row.get("summary_source"),
curation_score=value_or_default(row.get("curation_score"), 0),
info_density=value_or_default(row.get("info_density"), 50),
actionability=value_or_default(row.get("actionability"), 50),
diff --git a/backend/app/services/llm/_call_engine.py b/backend/app/services/llm/_call_engine.py
index 4dcef06f..0420f253 100644
--- a/backend/app/services/llm/_call_engine.py
+++ b/backend/app/services/llm/_call_engine.py
@@ -251,8 +251,11 @@ def _should_retry(exc: BaseException) -> bool:
覆盖 DeepSeek/GLM/智谱等通过 openai-compat 抛通用 APIError 的场景。
BadRequestError (400) 也不重试:内容过滤等确定性错误重试只会浪费时间。
+ 预算熔断拒绝(LlmBudgetExceededError)同理:窗口未恢复前重试必被再拒。
"""
- if isinstance(exc, RateLimitError) or _is_deterministic_request_error(exc):
+ from app.services.llm.budget_guard import LlmBudgetExceededError
+
+ if isinstance(exc, RateLimitError | LlmBudgetExceededError) or _is_deterministic_request_error(exc):
return False
return not _is_rate_limit_error(exc)
diff --git a/backend/app/services/llm/budget_guard.py b/backend/app/services/llm/budget_guard.py
new file mode 100644
index 00000000..e86c7d25
--- /dev/null
+++ b/backend/app/services/llm/budget_guard.py
@@ -0,0 +1,126 @@
+"""LLM 预算熔断闸。
+
+三窗(分钟/小时/日)预算保险丝:按 ``llm_call_logs`` 里实际 DONE 的调用
+计数,超限直接拒绝新调用(``LlmBudgetExceededError``),调用方走既有的
+失败退避/降级路径。与 ``_rate_limit`` 的语义区别:限流是「等一等再发」
+(保护对端配额),预算是「不再花下去」(保护自己的钱包)——拒绝而非等待。
+
+- 各窗口上限为 0 表示关闭该窗口;三窗全 0 时零开销直通。
+- 计数查询失败时 fail-open(保险丝自身故障不阻断主链路),记 warning。
+- 预算拒绝写一行 ``llm_call_logs`` 审计记录(status="BUDGET_REJECTED",
+ 不参与 DONE 计数),使拒绝在用量看板 ALL 视图可见。
+- 豁免场景(``LLM_BUDGET_EXEMPT_SCENES``,默认日报/周报/月报):低频
+ 核心承诺场景不检查预算——日报生成每天仅数次,不该被熔断硬失败
+ (daily_report 预算耗尽会走 ERROR 路径,违背核心承诺)。
+"""
+
+from __future__ import annotations
+
+import logging
+
+from app.core.config import settings
+from app.services.llm._rate_limit import record_llm_pool_circuit_event
+
+logger = logging.getLogger(__name__)
+
+
+class LlmBudgetExceededError(Exception):
+ """Rolling budget window exhausted; deterministic, must not be retried."""
+
+
+def _budget_limits() -> tuple[int, int, int] | None:
+ limits = (
+ settings.LLM_BUDGET_CALLS_PER_MINUTE,
+ settings.LLM_BUDGET_CALLS_PER_HOUR,
+ settings.LLM_BUDGET_CALLS_PER_DAY,
+ )
+ if not any(limits):
+ return None
+ return limits
+
+
+def _exempt_scenes() -> set[str]:
+ return {s.strip() for s in settings.LLM_BUDGET_EXEMPT_SCENES.split(",") if s.strip()}
+
+
+async def ensure_llm_budget(scene: str) -> None:
+ """Raise :class:`LlmBudgetExceededError` when a budget window is exhausted."""
+ if scene in _exempt_scenes():
+ return
+ limits = _budget_limits()
+ if limits is None:
+ return
+
+ from app.services.llm_usage import count_recent_llm_calls
+
+ try:
+ minute_calls, hour_calls, day_calls = await count_recent_llm_calls()
+ except Exception as exc: # noqa: BLE001 — fuse failure must not block the main path
+ logger.warning("LLM budget check skipped (count query failed): %s", exc)
+ return
+
+ for label, calls, limit in (
+ ("minute", minute_calls, limits[0]),
+ ("hour", hour_calls, limits[1]),
+ ("day", day_calls, limits[2]),
+ ):
+ if limit > 0 and calls >= limit:
+ logger.warning("LLM budget exceeded in %s window: %d/%d calls, scene=%s", label, calls, limit, scene)
+ record_llm_pool_circuit_event(None, scene, "budget_rejected")
+ await _record_budget_rejection(scene=scene, detail=f"{label} window: {calls}/{limit}")
+ raise LlmBudgetExceededError(
+ f"LLM budget exceeded in {label} window: {calls}/{limit} calls (scene={scene})"
+ )
+
+
+async def _record_budget_rejection(*, scene: str, detail: str) -> None:
+ """审计行:BUDGET_REJECTED 不进 DONE 计数,只在用量看板可见。"""
+ from app.services.llm_usage import record_llm_call_in_new_session
+
+ try:
+ await record_llm_call_in_new_session(
+ model=None,
+ request_model=None,
+ scene=scene,
+ status="BUDGET_REJECTED",
+ duration_ms=0,
+ error_message=detail,
+ )
+ except Exception as exc: # noqa: BLE001 — 审计写失败不能影响拒绝路径本身
+ logger.warning("LLM budget rejection audit log skipped: %s", exc)
+
+
+async def budget_headroom_ok(min_ratio: float = 0.3) -> bool:
+ """最紧启用窗口的剩余比例是否 ≥ ``min_ratio``,供降级内容补分析(requeue)门控。
+
+ 与 :func:`ensure_llm_budget` 的 fail-open 相反,这里 fail-closed:
+ requeue 是主动追加的负载,计数读不到时不应发起(宁可晚一轮补分析)。
+ 预算闸整体关闭(三窗全 0)时视为余量无限,恒 True。
+ """
+ limits = _budget_limits()
+ if limits is None:
+ return True
+
+ from app.services.llm_usage import count_recent_llm_calls
+
+ try:
+ calls = await count_recent_llm_calls()
+ except Exception as exc: # noqa: BLE001 — fail-closed:读不到计数就不补
+ logger.info("LLM budget headroom check skipped (count query failed): %s", exc)
+ return False
+
+ for label, used, limit in (
+ ("minute", calls[0], limits[0]),
+ ("hour", calls[1], limits[1]),
+ ("day", calls[2], limits[2]),
+ ):
+ if limit > 0 and (limit - used) / limit < min_ratio:
+ logger.info(
+ "LLM budget headroom below %.0f%% in %s window (%d/%d used), deferring requeue",
+ min_ratio * 100,
+ label,
+ used,
+ limit,
+ )
+ return False
+ return True
diff --git a/backend/app/services/llm/provider.py b/backend/app/services/llm/provider.py
index a8a0f672..49812f4b 100644
--- a/backend/app/services/llm/provider.py
+++ b/backend/app/services/llm/provider.py
@@ -152,6 +152,13 @@ async def call_llm_with_metadata(
if cached is not None:
return cached, {"cache_hit": True}
+ # Budget fuse: reject paid calls once a rolling budget window is exhausted.
+ # Checked after the response cache on purpose — cache hits are free and
+ # must keep serving while new paid calls are refused.
+ from app.services.llm.budget_guard import ensure_llm_budget
+
+ await ensure_llm_budget(scene)
+
try:
result = await _call_llm_with_metadata_inner(
messages,
@@ -172,12 +179,14 @@ async def call_llm_with_metadata(
return result
except Exception as exc:
# 输入或内容策略错误不反映模型可用性,不能污染全局熔断器。
+ from app.services.llm.budget_guard import LlmBudgetExceededError
from app.services.llm.circuit_breaker import CircuitOpenError
# 429 和本地候选冷却代表局部容量耗尽,由 per-model failover 管理;
# 把它们累计到路由熔断器会让一个配额不足的渠道阻断全部调用。
+ # 预算熔断同理:拒绝花钱不是故障,不得让预算恢复后仍被熔断。
if (
- not isinstance(exc, CircuitOpenError | LlmCapacityUnavailableError)
+ not isinstance(exc, CircuitOpenError | LlmCapacityUnavailableError | LlmBudgetExceededError)
and not _is_deterministic_request_error(exc)
and not _is_rate_limit_error(exc)
):
diff --git a/backend/app/services/llm_usage.py b/backend/app/services/llm_usage.py
index 037ef032..f1ffc9d8 100644
--- a/backend/app/services/llm_usage.py
+++ b/backend/app/services/llm_usage.py
@@ -323,3 +323,32 @@ async def _write():
except Exception as exc:
await db.rollback()
logger.warning("LLM usage log skipped: %s", exc)
+
+
+async def count_recent_llm_calls() -> tuple[int, int, int]:
+ """Count DONE LLM calls in the trailing 1m / 1h / 24h windows.
+
+ Single aggregate query backing the budget fuse (budget_guard); only
+ billable (status="DONE") calls count, rejected-by-fuse calls never
+ reach this table.
+ """
+ from datetime import UTC, datetime, timedelta
+
+ from sqlalchemy import case, func
+
+ from app.core.database import async_session
+
+ now = datetime.now(UTC)
+ async with async_session() as db:
+ result = await db.execute(
+ select(
+ func.count(case((LlmCallLog.created_at >= now - timedelta(minutes=1), 1))),
+ func.count(case((LlmCallLog.created_at >= now - timedelta(hours=1), 1))),
+ func.count(),
+ ).where(
+ LlmCallLog.status == "DONE",
+ LlmCallLog.created_at >= now - timedelta(hours=24),
+ )
+ )
+ minute_calls, hour_calls, day_calls = result.one()
+ return int(minute_calls), int(hour_calls), int(day_calls)
diff --git a/backend/app/services/scoring_engine.py b/backend/app/services/scoring_engine.py
index a1471bd3..c7a86084 100644
--- a/backend/app/services/scoring_engine.py
+++ b/backend/app/services/scoring_engine.py
@@ -89,15 +89,18 @@ class ScoringInput:
"source_weight_db", # Source.weight from DB (1-5)
# Feedback signal
"feedback_score", # user feedback signal (0+, default 0)
+ # Analysis provenance
+ "summary_source", # e.g. "local_fallback" for degraded local analysis (#90)
)
def __init__(self, **kwargs):
for slot in self.__slots__:
if slot == "feedback_score":
setattr(self, slot, kwargs.get(slot, 0))
- elif slot in ("published_at", "crawled_at"):
+ elif slot in ("published_at", "crawled_at", "summary_source"):
# datetime 字段默认 None(表示"没有"),不能是 0(int),
# 否则 _compute_time_decay 的 ensure_aware_utc(0) 会崩。
+ # summary_source 同理默认 None(表示真实 LLM 分析)。
setattr(self, slot, kwargs.get(slot))
else:
setattr(self, slot, kwargs.get(slot, 0))
@@ -442,7 +445,14 @@ def score_items(items: list[ScoringInput]) -> list[tuple[ScoreBreakdown, Scoring
if cfg["curation_mode"] == "percentile":
# Use final_score ranking: top (100 - percentile)% are selected
# e.g. curation_percentile=70 -> top 30% selected, bounded by base quality.
- final_scores = [bd.final_score for bd, _ in results]
+ # P70 门槛只由真实 LLM 评分决定:local_fallback 的确定性假分
+ # (curation 收敛 61-64 窄带)不得污染百分位门槛、挤掉真实内容。
+ # 降级项仍参与 selected 判定与展示(前端带「本地速览」标记),
+ # 由 requeue job 尽快换回真实分析(#90)。全候选皆降级时改用全量分数:
+ # 空列表会被 _compute_percentile_threshold 回退到全局默认阈值(55),
+ # 而降级内容的 final_score 普遍低于它,那样整批一条都选不出来。
+ real_scores = [bd.final_score for bd, item in results if item.summary_source != "local_fallback"]
+ final_scores = real_scores or [bd.final_score for bd, _ in results]
actual_threshold = _compute_percentile_threshold(
final_scores,
cfg["curation_percentile"],
diff --git a/backend/app/services/scoring_inputs.py b/backend/app/services/scoring_inputs.py
index 38be516c..4a13fbf5 100644
--- a/backend/app/services/scoring_inputs.py
+++ b/backend/app/services/scoring_inputs.py
@@ -44,6 +44,7 @@ def value_or_default(value: Any, default: float | int) -> float | int:
content_id=item.id,
title=item.title,
category=item.category,
+ summary_source=analysis.summary_source,
source_id=item.source_id,
source_name=item.source_name,
published_at=item.published_at,
diff --git a/backend/app/services/today_picks.py b/backend/app/services/today_picks.py
index 0a9c2172..d58aa937 100644
--- a/backend/app/services/today_picks.py
+++ b/backend/app/services/today_picks.py
@@ -6,7 +6,7 @@
import logging
from datetime import UTC, datetime, timedelta
-from sqlalchemy import or_, select
+from sqlalchemy import func, or_, select
from sqlalchemy.ext.asyncio import AsyncSession
from sqlalchemy.orm import selectinload
@@ -277,7 +277,9 @@ async def _build_today_picks_via_oltp(
selectinload(ContentItem.analyses),
selectinload(ContentItem.source),
)
- .where(ContentItem.crawled_at >= cutoff)
+ # 窗口按原文发布时间优先(published_at 为空的信源回退抓取时间),
+ # 旧文晚抓不进「今天」;与 DuckDB query_today_picks 的 COALESCE 口径对齐。
+ .where(func.coalesce(ContentItem.published_at, ContentItem.crawled_at) >= cutoff)
.order_by(ContentItem.crawled_at.desc())
)
if category:
@@ -429,6 +431,7 @@ def value_or_default(value, default):
content_id=row["id"],
title=row.get("title") or "",
category=row.get("category"),
+ summary_source=row.get("summary_source"),
source_id=row.get("source_id"),
source_name=row.get("source_name"),
published_at=row.get("published_at"),
@@ -478,6 +481,8 @@ def _row_to_content_payload(row: dict, breakdown: ScoreBreakdown) -> dict:
"short_video_plan": _decode_json_value(row.get("short_video_plan")),
"risk_notes": _decode_json_value(row.get("risk_notes")),
"curation_score": row.get("curation_score") or 0,
+ # 分析来源:'local_fallback' 时前端显示「本地速览·待 AI 复核」标记(#90)
+ "summary_source": row.get("summary_source"),
"tags": analysis_tags,
"recommendation": _clean_optional_text(row.get("recommendation")),
"info_density": row.get("info_density") or 0,
diff --git a/backend/tests/test_analysis_requeue.py b/backend/tests/test_analysis_requeue.py
new file mode 100644
index 00000000..d5e1b549
--- /dev/null
+++ b/backend/tests/test_analysis_requeue.py
@@ -0,0 +1,354 @@
+"""降级内容补分析(requeue)与预算豁免/余量门控的回归测试(#90)。
+
+覆盖:
+- repo.reset_local_fallback_for_reanalysis:资格(最新分析是 local_fallback)、
+ cooldown、skip_analysis、max_age、limit 最旧优先、reset 字段语义;
+- requeue_local_fallback_content:预算余量门控;
+- budget_guard.budget_headroom_ok:全关、余量足、不足、fail-closed;
+- 预算豁免场景(daily_report 等)超限仍放行。
+"""
+
+from __future__ import annotations
+
+from datetime import UTC, datetime, timedelta
+
+import pytest
+from sqlalchemy import select
+from sqlalchemy.ext.asyncio import async_sessionmaker, create_async_engine
+
+from app.core.config import settings
+from app.core.database import Base
+from app.models.analysis import AiAnalysis
+from app.models.content import ContentItem, ContentStatus
+from app.models.source import Source, SourceStatus, SourceType
+from app.repositories.content_repo import ContentRepo
+from app.services.analysis_requeue import requeue_local_fallback_content
+from app.services.llm.budget_guard import (
+ LlmBudgetExceededError,
+ budget_headroom_ok,
+ ensure_llm_budget,
+)
+
+
+def _set_budget(monkeypatch, minute: int, hour: int, day: int) -> None:
+ monkeypatch.setattr(settings, "LLM_BUDGET_CALLS_PER_MINUTE", minute)
+ monkeypatch.setattr(settings, "LLM_BUDGET_CALLS_PER_HOUR", hour)
+ monkeypatch.setattr(settings, "LLM_BUDGET_CALLS_PER_DAY", day)
+
+
+def _content(
+ item_id: int,
+ *,
+ now: datetime,
+ updated_at: datetime | None = None,
+ created_at: datetime | None = None,
+ skip_analysis: bool = False,
+ status: ContentStatus = ContentStatus.ANALYZED,
+) -> ContentItem:
+ return ContentItem(
+ id=item_id,
+ title=f"降级样本 {item_id}",
+ url=f"https://example.com/fallback-{item_id}",
+ source_id=1,
+ source_name="降级测试信源",
+ source_type="RSS",
+ category="AI",
+ status=status,
+ skip_analysis=skip_analysis,
+ crawled_at=now,
+ published_at=now,
+ created_at=created_at or now,
+ updated_at=updated_at or now - timedelta(minutes=120),
+ )
+
+
+async def _seed(db, now: datetime) -> None:
+ db.add(
+ Source(
+ id=1,
+ name="降级测试信源",
+ source_type=SourceType.RSS,
+ url="https://example.com/fallback.xml",
+ category="AI",
+ status=SourceStatus.ACTIVE,
+ enabled=True,
+ weight=3,
+ )
+ )
+ # id=1:最新分析是 local_fallback,超过 cooldown —— 应被回收
+ db.add(_content(1, now=now))
+ # id=2:最新分析是真实 LLM(fallback 之后又真实分析过)—— 不动
+ db.add(_content(2, now=now))
+ # id=3:fallback 但 cooldown 内(updated_at 刚刷新)—— 不动
+ db.add(_content(3, now=now, updated_at=now - timedelta(minutes=5)))
+ # id=4:fallback 但显式 skip_analysis —— 不动
+ db.add(_content(4, now=now, skip_analysis=True))
+ # id=5:fallback 但超过 max_age_days —— 不动
+ db.add(_content(5, now=now, created_at=now - timedelta(days=10)))
+
+ analyses = [
+ # id=1:一条真实分析(旧)+ 一条 fallback(新)→ 最新是 fallback
+ AiAnalysis(id=101, content_id=1, curation_score=80, created_at=now - timedelta(hours=3)),
+ AiAnalysis(
+ id=102,
+ content_id=1,
+ curation_score=62,
+ summary_source="local_fallback",
+ created_at=now - timedelta(hours=2),
+ ),
+ # id=2:fallback(旧)+ 真实分析(新)→ 最新是真实
+ AiAnalysis(
+ id=103,
+ content_id=2,
+ curation_score=62,
+ summary_source="local_fallback",
+ created_at=now - timedelta(hours=3),
+ ),
+ AiAnalysis(id=104, content_id=2, curation_score=85, created_at=now - timedelta(hours=1)),
+ # id=3/4/5:最新分析均为 fallback
+ AiAnalysis(
+ id=105,
+ content_id=3,
+ curation_score=62,
+ summary_source="local_fallback",
+ created_at=now - timedelta(minutes=10),
+ ),
+ AiAnalysis(
+ id=106,
+ content_id=4,
+ curation_score=62,
+ summary_source="local_fallback",
+ created_at=now - timedelta(hours=2),
+ ),
+ AiAnalysis(
+ id=107,
+ content_id=5,
+ curation_score=62,
+ summary_source="local_fallback",
+ created_at=now - timedelta(days=9),
+ ),
+ ]
+ db.add_all(analyses)
+ await db.commit()
+
+
+@pytest.mark.asyncio
+async def test_reset_local_fallback_eligibility_and_fields():
+ """只有「最新分析是 fallback 且过冷却」的内容被回收,reset 字段语义正确。"""
+ engine = create_async_engine("sqlite+aiosqlite:///:memory:")
+ session_factory = async_sessionmaker(engine, expire_on_commit=False)
+ async with engine.begin() as conn:
+ await conn.run_sync(Base.metadata.create_all)
+
+ now = datetime.now(UTC)
+ async with session_factory() as db:
+ await _seed(db, now)
+ repo = ContentRepo(db)
+ count = await repo.reset_local_fallback_for_reanalysis(limit=10, cooldown_minutes=60, max_age_days=7)
+ assert count == 1, "只有 id=1 同时满足最新分析为 fallback、过冷却、未跳过、未过期"
+
+ item1 = await db.get(ContentItem, 1)
+ assert item1 is not None
+ assert item1.status == ContentStatus.PENDING
+ assert item1.analysis_attempts == 0
+ assert item1.analysis_next_retry_at is None
+ assert item1.analysis_claim_token is None
+ assert item1.analysis_lease_expires_at is None
+
+ for unchanged_id in (2, 3, 4, 5):
+ item = await db.get(ContentItem, unchanged_id)
+ assert item is not None
+ assert item.status == ContentStatus.ANALYZED, f"id={unchanged_id} 不应被回收"
+
+ await engine.dispose()
+
+
+@pytest.mark.asyncio
+async def test_reset_respects_limit_oldest_first():
+ """limit 限速且最旧优先。"""
+ engine = create_async_engine("sqlite+aiosqlite:///:memory:")
+ session_factory = async_sessionmaker(engine, expire_on_commit=False)
+ async with engine.begin() as conn:
+ await conn.run_sync(Base.metadata.create_all)
+
+ now = datetime.now(UTC)
+ async with session_factory() as db:
+ db.add(
+ Source(
+ id=1,
+ name="信源",
+ source_type=SourceType.RSS,
+ url="https://example.com/x.xml",
+ category="AI",
+ status=SourceStatus.ACTIVE,
+ enabled=True,
+ weight=3,
+ )
+ )
+ # 三条 fallback 内容,created_at 递减(id=10 最旧)
+ for item_id, age_hours in ((10, 6), (11, 4), (12, 2)):
+ db.add(
+ _content(
+ item_id,
+ now=now,
+ created_at=now - timedelta(hours=age_hours),
+ updated_at=now - timedelta(hours=age_hours),
+ )
+ )
+ db.add(
+ AiAnalysis(
+ id=200 + item_id,
+ content_id=item_id,
+ curation_score=62,
+ summary_source="local_fallback",
+ created_at=now - timedelta(hours=age_hours),
+ )
+ )
+ await db.commit()
+
+ repo = ContentRepo(db)
+ count = await repo.reset_local_fallback_for_reanalysis(limit=1, cooldown_minutes=60)
+ assert count == 1
+ pending_ids = [
+ int(row[0])
+ for row in (
+ await db.execute(select(ContentItem.id).where(ContentItem.status == ContentStatus.PENDING))
+ ).all()
+ ]
+ assert pending_ids == [10], "limit=1 时应回收最旧的 id=10"
+
+ await engine.dispose()
+
+
+@pytest.mark.asyncio
+async def test_requeue_service_gated_by_budget_headroom(monkeypatch):
+ """预算余量不足时不发起回收(不触碰 repo)。"""
+ from unittest.mock import AsyncMock
+
+ reset_mock = AsyncMock(return_value=5)
+ monkeypatch.setattr(ContentRepo, "reset_local_fallback_for_reanalysis", reset_mock)
+
+ async def _no_headroom(min_ratio: float = 0.3) -> bool:
+ return False
+
+ monkeypatch.setattr("app.services.analysis_requeue.budget_headroom_ok", _no_headroom)
+
+ count = await requeue_local_fallback_content(object()) # db 不应被触碰
+ assert count == 0
+ reset_mock.assert_not_awaited()
+
+
+@pytest.mark.asyncio
+async def test_budget_headroom_semantics(monkeypatch):
+ # 全关(三窗全 0)→ 余量无限
+ _set_budget(monkeypatch, 0, 0, 0)
+
+ async def _boom():
+ raise AssertionError("budget disabled must not query counts")
+
+ monkeypatch.setattr("app.services.llm_usage.count_recent_llm_calls", _boom)
+ assert await budget_headroom_ok() is True
+
+ # 余量充足
+ _set_budget(monkeypatch, 100, 2000, 20000)
+
+ async def _counts():
+ return 10, 500, 10000
+
+ monkeypatch.setattr("app.services.llm_usage.count_recent_llm_calls", _counts)
+ assert await budget_headroom_ok() is True
+
+ # 最紧窗口(日窗 10000/20000 = 50% 已用,余量 50% ≥30%);改为 18000(余量 10%)→ False
+ async def _tight():
+ return 5, 100, 18000
+
+ monkeypatch.setattr("app.services.llm_usage.count_recent_llm_calls", _tight)
+ assert await budget_headroom_ok() is False
+
+ # fail-closed:计数查询失败 → False(requeue 不发起)
+ async def _db_error():
+ raise RuntimeError("db unavailable")
+
+ monkeypatch.setattr("app.services.llm_usage.count_recent_llm_calls", _db_error)
+ assert await budget_headroom_ok() is False
+
+
+@pytest.mark.asyncio
+async def test_exempt_scenes_bypass_budget(monkeypatch):
+ """日报/周报/月报等豁免场景超限仍放行;非豁免场景照常拒绝。"""
+ _set_budget(monkeypatch, minute=1, hour=0, day=0)
+
+ async def _counts():
+ return 1, 0, 0 # 分钟窗已耗尽
+
+ monkeypatch.setattr("app.services.llm_usage.count_recent_llm_calls", _counts)
+
+ for scene in ("daily_report", "weekly_digest", "monthly_digest"):
+ await ensure_llm_budget(scene) # 豁免:不抛
+
+ with pytest.raises(LlmBudgetExceededError):
+ await ensure_llm_budget("content_analysis")
+
+
+@pytest.mark.asyncio
+async def test_reset_skips_content_with_repeated_fallbacks():
+ """已有 ≥2 条 local_fallback 分析的内容不再回收(防确定性降级无限循环)。"""
+ engine = create_async_engine("sqlite+aiosqlite:///:memory:")
+ session_factory = async_sessionmaker(engine, expire_on_commit=False)
+ async with engine.begin() as conn:
+ await conn.run_sync(Base.metadata.create_all)
+
+ now = datetime.now(UTC)
+ async with session_factory() as db:
+ db.add(
+ Source(
+ id=1,
+ name="信源",
+ source_type=SourceType.RSS,
+ url="https://example.com/y.xml",
+ category="AI",
+ status=SourceStatus.ACTIVE,
+ enabled=True,
+ weight=3,
+ )
+ )
+ # id=20:两条 fallback 分析(requeue 重降级过一次)—— 不再回收
+ db.add(_content(20, now=now))
+ # id=21:一条 fallback —— 正常回收(对照组)
+ db.add(_content(21, now=now))
+ db.add_all(
+ [
+ AiAnalysis(
+ id=301,
+ content_id=20,
+ curation_score=62,
+ summary_source="local_fallback",
+ created_at=now - timedelta(hours=4),
+ ),
+ AiAnalysis(
+ id=302,
+ content_id=20,
+ curation_score=62,
+ summary_source="local_fallback",
+ created_at=now - timedelta(hours=2),
+ ),
+ AiAnalysis(
+ id=303,
+ content_id=21,
+ curation_score=62,
+ summary_source="local_fallback",
+ created_at=now - timedelta(hours=2),
+ ),
+ ]
+ )
+ await db.commit()
+
+ repo = ContentRepo(db)
+ count = await repo.reset_local_fallback_for_reanalysis(limit=10)
+ assert count == 1
+ item20 = await db.get(ContentItem, 20)
+ item21 = await db.get(ContentItem, 21)
+ assert item20 is not None and item20.status == ContentStatus.ANALYZED
+ assert item21 is not None and item21.status == ContentStatus.PENDING
+
+ await engine.dispose()
diff --git a/backend/tests/test_duckdb_service.py b/backend/tests/test_duckdb_service.py
index 8fa4e2cb..89378075 100644
--- a/backend/tests/test_duckdb_service.py
+++ b/backend/tests/test_duckdb_service.py
@@ -826,7 +826,8 @@ def test_digest_content_query_uses_latest_analysis_and_feedback_order(monkeypatc
category VARCHAR,
source_name VARCHAR,
platform VARCHAR,
- crawled_at TIMESTAMP
+ crawled_at TIMESTAMP,
+ published_at TIMESTAMP
)
""")
conn.execute("""
@@ -847,7 +848,8 @@ def test_digest_content_query_uses_latest_analysis_and_feedback_order(monkeypatc
tags VARCHAR,
recommendation VARCHAR,
recommended_reason VARCHAR,
- created_at TIMESTAMP
+ created_at TIMESTAMP,
+ summary_source VARCHAR
)
""")
conn.execute("""
@@ -873,23 +875,23 @@ def test_digest_content_query_uses_latest_analysis_and_feedback_order(monkeypatc
now = datetime.now(UTC).replace(tzinfo=None).replace(tzinfo=None)
conn.execute("INSERT INTO oltp_db.sources VALUES (10, 4)")
conn.execute(
- "INSERT INTO oltp_db.content_items VALUES (1, 10, '反馈后的最新分析', 'https://example.com/1', 'AI', '测试信源', 'rss', ?)",
- [now],
+ "INSERT INTO oltp_db.content_items VALUES (1, 10, '反馈后的最新分析', 'https://example.com/1', 'AI', '测试信源', 'rss', ?, ?)",
+ [now, now],
)
conn.execute(
- "INSERT INTO oltp_db.content_items VALUES (2, NULL, '无反馈样本', 'https://example.com/2', 'AI', '测试信源', 'rss', ?)",
- [now],
+ "INSERT INTO oltp_db.content_items VALUES (2, NULL, '无反馈样本', 'https://example.com/2', 'AI', '测试信源', 'rss', ?, ?)",
+ [now, now],
)
conn.execute(
- "INSERT INTO oltp_db.ai_analyses VALUES (1, 1, '旧摘要', 95, 95, 95, 95, 95, 10, 99, 95, 95, 95, '[\"旧\"]', '旧推荐', '旧理由', ?)",
+ "INSERT INTO oltp_db.ai_analyses VALUES (1, 1, '旧摘要', 95, 95, 95, 95, 95, 10, 99, 95, 95, 95, '[\"旧\"]', '旧推荐', '旧理由', ?, NULL)",
[now - timedelta(hours=2)],
)
conn.execute(
- "INSERT INTO oltp_db.ai_analyses VALUES (2, 1, '新摘要', 70, 70, 70, 65, 80, 10, 70, 82, 81, 72, '[\"新\"]', '新推荐', '新理由', ?)",
+ "INSERT INTO oltp_db.ai_analyses VALUES (2, 1, '新摘要', 70, 70, 70, 65, 80, 10, 70, 82, 81, 72, '[\"新\"]', '新推荐', '新理由', ?, NULL)",
[now - timedelta(hours=1)],
)
conn.execute(
- "INSERT INTO oltp_db.ai_analyses VALUES (3, 2, '无反馈摘要', 72, 72, 72, 60, 75, 10, 72, 74, 73, 50, '[\"AI\"]', '推荐', '理由', ?)",
+ "INSERT INTO oltp_db.ai_analyses VALUES (3, 2, '无反馈摘要', 72, 72, 72, 60, 75, 10, 72, 74, 73, 50, '[\"AI\"]', '推荐', '理由', ?, NULL)",
[now],
)
conn.execute("INSERT INTO oltp_db.user_feedback VALUES (1, 1, 1, 20.0, ?)", [now])
@@ -930,7 +932,8 @@ def test_daily_report_content_query_uses_latest_analysis_only(monkeypatch):
url VARCHAR,
category VARCHAR,
source_name VARCHAR,
- crawled_at TIMESTAMP
+ crawled_at TIMESTAMP,
+ published_at TIMESTAMP
)
""")
conn.execute("""
@@ -951,12 +954,12 @@ def test_daily_report_content_query_uses_latest_analysis_only(monkeypatch):
now = datetime.now(UTC).replace(tzinfo=None)
conn.execute(
- "INSERT INTO oltp_db.content_items VALUES (1, '多次分析样本', 'https://example.com/1', 'AI', '测试信源', ?)",
- [now],
+ "INSERT INTO oltp_db.content_items VALUES (1, '多次分析样本', 'https://example.com/1', 'AI', '测试信源', ?, ?)",
+ [now, now],
)
conn.execute(
- "INSERT INTO oltp_db.content_items VALUES (2, '普通样本', 'https://example.com/2', 'AI', '测试信源', ?)",
- [now],
+ "INSERT INTO oltp_db.content_items VALUES (2, '普通样本', 'https://example.com/2', 'AI', '测试信源', ?, ?)",
+ [now, now],
)
conn.execute(
"INSERT INTO oltp_db.ai_analyses VALUES (1, 1, '旧日报摘要', 99, 99, 99, 10, 99, '旧理由', ?)",
@@ -996,7 +999,8 @@ def test_digest_content_queries_exclude_ignored_content(monkeypatch):
category VARCHAR,
source_name VARCHAR,
platform VARCHAR,
- crawled_at TIMESTAMP
+ crawled_at TIMESTAMP,
+ published_at TIMESTAMP
)
""")
conn.execute("""
@@ -1017,7 +1021,8 @@ def test_digest_content_queries_exclude_ignored_content(monkeypatch):
tags VARCHAR,
recommendation VARCHAR,
recommended_reason VARCHAR,
- created_at TIMESTAMP
+ created_at TIMESTAMP,
+ summary_source VARCHAR
)
""")
conn.execute("""
@@ -1040,13 +1045,13 @@ def test_digest_content_queries_exclude_ignored_content(monkeypatch):
now = datetime.now(UTC).replace(tzinfo=None)
for content_id, title in ((1, "已忽略素材"), (2, "保留素材")):
conn.execute(
- "INSERT INTO oltp_db.content_items VALUES (?, NULL, ?, ?, 'AI', '测试信源', 'rss', ?)",
- [content_id, title, f"https://example.com/{content_id}", now],
+ "INSERT INTO oltp_db.content_items VALUES (?, NULL, ?, ?, 'AI', '测试信源', 'rss', ?, ?)",
+ [content_id, title, f"https://example.com/{content_id}", now, now],
)
conn.execute(
"""
INSERT INTO oltp_db.ai_analyses VALUES (
- ?, ?, '摘要', 90, 90, 90, 90, 90, 10, 90, 90, 90, 90, '["AI"]', '推荐', '理由', ?
+ ?, ?, '摘要', 90, 90, 90, 90, 90, 10, 90, 90, 90, 90, '["AI"]', '推荐', '理由', ?, NULL
)
""",
[content_id, content_id, now],
diff --git a/backend/tests/test_llm_budget_guard.py b/backend/tests/test_llm_budget_guard.py
new file mode 100644
index 00000000..e8c0466f
--- /dev/null
+++ b/backend/tests/test_llm_budget_guard.py
@@ -0,0 +1,205 @@
+"""LLM 预算熔断闸(budget_guard)的行为回归测试。
+
+覆盖三窗语义:
+- 超限拒绝(分钟/日窗各自生效),拒绝事件可观测(budget_rejected 事件 + warning);
+- 三窗全 0 时零开销直通(不触发计数查询);
+- 计数查询失败 fail-open(保险丝自身故障不阻断主链路);
+- provider 集成:预算拒绝不进入 failover/真实调用、不污染路由熔断器、
+ 响应缓存命中(免费)在预算耗尽时仍放行。
+"""
+
+from __future__ import annotations
+
+import pytest
+
+from app.core.config import settings
+from app.services.llm import provider
+from app.services.llm._call_engine import _should_retry
+from app.services.llm.budget_guard import (
+ LlmBudgetExceededError,
+ _budget_limits,
+ ensure_llm_budget,
+)
+from app.services.llm.circuit_breaker import (
+ get_llm_circuit_breaker,
+ reset_llm_circuit_breakers,
+)
+from app.services.llm.response_cache import get_llm_cache
+
+
+def _set_budget(monkeypatch, minute: int, hour: int, day: int) -> None:
+ monkeypatch.setattr(settings, "LLM_BUDGET_CALLS_PER_MINUTE", minute)
+ monkeypatch.setattr(settings, "LLM_BUDGET_CALLS_PER_HOUR", hour)
+ monkeypatch.setattr(settings, "LLM_BUDGET_CALLS_PER_DAY", day)
+
+
+@pytest.fixture(autouse=True)
+def _isolate_route_state():
+ provider._failover.reset()
+ reset_llm_circuit_breakers()
+ get_llm_cache().clear()
+ yield
+ provider._failover.reset()
+ reset_llm_circuit_breakers()
+ get_llm_cache().clear()
+
+
+@pytest.mark.asyncio
+async def test_all_zero_budget_skips_count_query(monkeypatch):
+ """三窗全 0 = 预算闸关闭:不应触碰计数查询。"""
+ _set_budget(monkeypatch, 0, 0, 0)
+ assert _budget_limits() is None
+
+ async def _boom():
+ raise AssertionError("budget disabled must not query call counts")
+
+ monkeypatch.setattr("app.services.llm_usage.count_recent_llm_calls", _boom)
+ await ensure_llm_budget(scene="test")
+
+
+@pytest.mark.asyncio
+async def test_minute_window_rejects_and_emits_event(monkeypatch):
+ _set_budget(monkeypatch, minute=2, hour=0, day=0)
+
+ async def _counts():
+ return 2, 5, 10 # 分钟窗已到 2/2
+
+ monkeypatch.setattr("app.services.llm_usage.count_recent_llm_calls", _counts)
+ events: list[tuple] = []
+ monkeypatch.setattr(
+ "app.services.llm.budget_guard.record_llm_pool_circuit_event",
+ lambda *args: events.append(args),
+ )
+ audits: list[dict] = []
+
+ async def _fake_audit(**kwargs):
+ audits.append(kwargs)
+
+ monkeypatch.setattr("app.services.llm.budget_guard._record_budget_rejection", _fake_audit)
+
+ with pytest.raises(LlmBudgetExceededError, match="minute"):
+ await ensure_llm_budget(scene="analysis")
+ assert events and events[0][2] == "budget_rejected"
+ assert audits and audits[0]["scene"] == "analysis"
+ assert "minute window: 2/2" in audits[0]["detail"]
+
+
+@pytest.mark.asyncio
+async def test_day_window_rejects(monkeypatch):
+ _set_budget(monkeypatch, minute=0, hour=0, day=100)
+
+ async def _counts():
+ return 1, 50, 100 # 日窗已到 100/100
+
+ monkeypatch.setattr("app.services.llm_usage.count_recent_llm_calls", _counts)
+
+ async def _fake_audit(**kwargs):
+ return None
+
+ monkeypatch.setattr("app.services.llm.budget_guard._record_budget_rejection", _fake_audit)
+ with pytest.raises(LlmBudgetExceededError, match="day"):
+ await ensure_llm_budget(scene="analysis")
+
+
+@pytest.mark.asyncio
+async def test_hour_window_rejects(monkeypatch):
+ _set_budget(monkeypatch, minute=0, hour=500, day=0)
+
+ async def _counts():
+ return 10, 500, 600 # 小时窗已到 500/500
+
+ monkeypatch.setattr("app.services.llm_usage.count_recent_llm_calls", _counts)
+
+ async def _fake_audit(**kwargs):
+ return None
+
+ monkeypatch.setattr("app.services.llm.budget_guard._record_budget_rejection", _fake_audit)
+ with pytest.raises(LlmBudgetExceededError, match="hour"):
+ await ensure_llm_budget(scene="analysis")
+
+
+@pytest.mark.asyncio
+async def test_within_budget_passes(monkeypatch):
+ _set_budget(monkeypatch, minute=10, hour=100, day=1000)
+
+ async def _counts():
+ return 3, 40, 500
+
+ monkeypatch.setattr("app.services.llm_usage.count_recent_llm_calls", _counts)
+ await ensure_llm_budget(scene="analysis") # 不抛即通过
+
+
+@pytest.mark.asyncio
+async def test_count_query_failure_fails_open(monkeypatch, caplog):
+ _set_budget(monkeypatch, minute=10, hour=100, day=1000)
+
+ async def _boom():
+ raise RuntimeError("db unavailable")
+
+ monkeypatch.setattr("app.services.llm_usage.count_recent_llm_calls", _boom)
+ with caplog.at_level("WARNING"):
+ await ensure_llm_budget(scene="analysis") # fail-open,不抛
+ assert any("budget check skipped" in rec.message for rec in caplog.records)
+
+
+def test_should_retry_rejects_budget_error():
+ assert _should_retry(LlmBudgetExceededError("minute window")) is False
+
+
+@pytest.mark.asyncio
+async def test_provider_rejects_before_failover_and_keeps_breaker_clean(monkeypatch):
+ """预算拒绝:不进入真实调用/failover,也不计入路由熔断器失败。"""
+ _set_budget(monkeypatch, minute=1, hour=0, day=0)
+
+ async def _counts():
+ return 1, 0, 0
+
+ monkeypatch.setattr("app.services.llm_usage.count_recent_llm_calls", _counts)
+
+ async def _forbidden(*args, **kwargs):
+ raise AssertionError("budget-exceeded call must not reach the failover layer")
+
+ monkeypatch.setattr(provider, "_call_llm_with_metadata_inner", _forbidden)
+
+ breaker = get_llm_circuit_breaker("budget_test")
+ failures_before = breaker.status()["failure_count"]
+
+ with pytest.raises(LlmBudgetExceededError):
+ await provider.call_llm_with_metadata(
+ [{"role": "user", "content": "预算耗尽时的调用"}],
+ routing_group="budget_test",
+ scene="test",
+ )
+ assert breaker.status()["failure_count"] == failures_before, "预算拒绝不是故障,不得污染熔断器"
+
+
+@pytest.mark.asyncio
+async def test_cache_hit_still_served_when_budget_exhausted(monkeypatch):
+ """预算耗尽时缓存命中(免费)必须继续放行。"""
+ _set_budget(monkeypatch, minute=100, hour=0, day=0)
+
+ async def _inner(messages, temperature, max_tokens, scene, routing_group, response_format=None):
+ return "cached-answer", {"model": "m1"}
+
+ async def _zero_counts():
+ return 0, 0, 0
+
+ monkeypatch.setattr("app.services.llm_usage.count_recent_llm_calls", _zero_counts)
+ monkeypatch.setattr(provider, "_call_llm_with_metadata_inner", _inner)
+
+ messages = [{"role": "user", "content": "会被缓存的问题"}]
+ text, meta = await provider.call_llm_with_metadata(messages, routing_group="budget_cache", scene="test")
+ assert text == "cached-answer"
+ assert meta.get("cache_hit") is not True
+
+ # 预算降到已耗尽,同参调用必须仍能从缓存拿到答案
+ _set_budget(monkeypatch, minute=1, hour=0, day=0)
+
+ async def _counts():
+ return 1, 0, 0
+
+ monkeypatch.setattr("app.services.llm_usage.count_recent_llm_calls", _counts)
+
+ text2, meta2 = await provider.call_llm_with_metadata(messages, routing_group="budget_cache", scene="test")
+ assert text2 == "cached-answer"
+ assert meta2.get("cache_hit") is True
diff --git a/backend/tests/test_scoring_engine.py b/backend/tests/test_scoring_engine.py
index 83ef106a..82a14e37 100644
--- a/backend/tests/test_scoring_engine.py
+++ b/backend/tests/test_scoring_engine.py
@@ -1,6 +1,11 @@
from datetime import UTC, datetime, timedelta
-from app.services.scoring_engine import CONFIG, ScoringInput, score_items
+from app.services.scoring_engine import (
+ CONFIG,
+ ScoringInput,
+ _compute_percentile_threshold,
+ score_items,
+)
_NOW = datetime(2026, 1, 1, 12, 0, 0)
@@ -210,3 +215,73 @@ def test_diversity_promotes_other_source_over_same_source_run():
# 第一名仍是分数最高的同源第一条;同源第二条因惩罚被其它来源反超
assert final_ids[0] == 1
assert final_ids.index(99) < final_ids.index(2)
+
+
+def test_p70_threshold_excludes_local_fallback_fake_scores():
+ """#90:local_fallback 的确定性假分不得污染 P70 门槛。
+
+ 三组对照(自证明区分度——守门断言保证若删掉排除逻辑本测试必红):
+ - real_only:3 条真实分(90/80/70),门槛 T1 只由真实分决定;
+ - mixed:再混入 2 条假分(85/75,标记 local_fallback)——真实项门槛
+ 仍应为 T1(排除生效);
+ - full:同样的 2 条 85/75 但标记为真实分析——门槛 T2 应不同于 T1,
+ 证明该分数组合确实会移动门槛(若排除逻辑失效,mixed 会退化成
+ full 的门槛,测试转红)。
+ """
+
+ def _scored(content_id: int, curation: int, **extra) -> ScoringInput:
+ return _item(content_id, curation_score=curation, **extra)
+
+ real_items = [_scored(1, 90), _scored(2, 80), _scored(3, 70)]
+ fallback_items = [
+ _scored(98, 85, summary_source="local_fallback"),
+ _scored(99, 75, summary_source="local_fallback"),
+ ]
+ same_scores_as_real = [_scored(98, 85), _scored(99, 75)] # 不带标记
+
+ baseline = score_items(list(real_items))
+ mixed = score_items(real_items + fallback_items)
+ full = score_items(real_items + same_scores_as_real)
+
+ t_baseline = {bd.threshold_used for bd, _ in baseline}
+ t_mixed_real = {bd.threshold_used for bd, item in mixed if item.content_id in (1, 2, 3)}
+ t_full = {bd.threshold_used for bd, _ in full}
+
+ assert len(t_baseline) == 1 and len(t_full) == 1
+ assert t_mixed_real == t_baseline, "混入 fallback 假分后真实项门槛不应改变"
+ assert t_full != t_baseline, "守门断言:85/75 混入真实批应移动门槛——本断言失败说明对照构造失去区分度"
+
+ # fallback 项自身仍参与判定与展示(不静默消失),且用同一门槛判定
+ fallback_results = {item.content_id: bd for bd, item in mixed if item.content_id in (98, 99)}
+ assert len(fallback_results) == 2
+ assert {bd.threshold_used for bd in fallback_results.values()} == t_baseline
+
+
+def test_all_fallback_batch_falls_back_to_full_scores():
+ """全候选皆降级时,门槛仍由本批实际分数决定,不退回全局默认阈值。
+
+ 这个分支不是防御性冗余:真实分为空时若直接把空列表交给
+ ``_compute_percentile_threshold``,它会返回 ``CONFIG["curation_threshold"]``
+ (55),而降级内容的 final_score 普遍低于该值,结果是整批一条都选不出来。
+ """
+ fallback_items = [_item(i, summary_source="local_fallback", curation_score=62) for i in range(1, 5)]
+ scored = score_items(fallback_items)
+ assert len(scored) == 4
+
+ thresholds = {bd.threshold_used for bd, _ in scored}
+ assert len(thresholds) == 1
+ threshold = thresholds.pop()
+
+ own_scores = [bd.final_score for bd, _ in scored]
+ assert threshold == _compute_percentile_threshold(own_scores, 70), "门槛应由本批实际分数决定"
+
+ # 守门断言:门槛退回全局默认时本批会全部低于阈值、页面一条都不出。
+ # 删掉 score_items 里的 `or [bd.final_score ...]` 回退分支,本测试必红。
+ assert threshold != CONFIG["curation_threshold"]
+ assert any(bd.selected for bd, _ in scored), "门槛来自本批分数时至少应选出一条"
+
+
+def test_scoring_input_defaults_summary_source_none():
+ """未传 summary_source 的旧构造路径默认 None(表示真实 LLM 分析)。"""
+ item = _item(1)
+ assert item.summary_source is None
diff --git a/backend/tests/test_scoring_flow_repo.py b/backend/tests/test_scoring_flow_repo.py
index b52656c3..c0f14c78 100644
--- a/backend/tests/test_scoring_flow_repo.py
+++ b/backend/tests/test_scoring_flow_repo.py
@@ -654,3 +654,110 @@ async def test_scoring_flow_counts_respect_visible_user_id():
assert count_all == 3, f"全局应见 3 条,实际 {count_all}"
await engine.dispose()
+
+
+def _window_test_source(now: datetime) -> Source:
+ return Source(
+ id=1,
+ name="窗口测试信源",
+ source_type=SourceType.RSS,
+ url="https://example.com/window.xml",
+ category="AI",
+ status=SourceStatus.ACTIVE,
+ enabled=True,
+ weight=3,
+ )
+
+
+def _analyzed_item(
+ item_id: int,
+ *,
+ crawled_at: datetime,
+ published_at: datetime | None,
+) -> ContentItem:
+ return ContentItem(
+ id=item_id,
+ title=f"窗口样本 {item_id}",
+ url=f"https://example.com/window-{item_id}",
+ source_id=1,
+ source_name="窗口测试信源",
+ source_type="RSS",
+ category="AI",
+ status=ContentStatus.ANALYZED,
+ crawled_at=crawled_at,
+ published_at=published_at,
+ )
+
+
+async def _seed_window_items(db, now: datetime) -> None:
+ """三条内容:旧文晚抓 / 无发布时间 / 正常新文,各带一条低风险分析。"""
+ items = [
+ # 旧文晚抓:3 天前发布、刚刚才发现 —— 不应进「今天」的候选
+ _analyzed_item(1, crawled_at=now, published_at=now - timedelta(days=3)),
+ # 无发布时间的信源:回退按抓取时间计(2 小时前抓到)
+ _analyzed_item(2, crawled_at=now - timedelta(hours=2), published_at=None),
+ # 正常新文:半小时前发布、刚刚抓到(发布时间在报告窗口内)
+ _analyzed_item(3, crawled_at=now, published_at=now - timedelta(minutes=30)),
+ ]
+ db.add_all(items)
+ db.add_all(
+ [
+ AiAnalysis(
+ id=item_id,
+ content_id=item_id,
+ curation_score=80,
+ risk_score=10,
+ created_at=now,
+ )
+ for item_id in (1, 2, 3)
+ ]
+ )
+ await db.commit()
+
+
+@pytest.mark.asyncio
+async def test_today_picks_window_prefers_published_at():
+ """旧文晚抓不进今日候选;published_at 缺失回退 crawled_at。"""
+ engine = create_async_engine("sqlite+aiosqlite:///:memory:")
+ session_factory = async_sessionmaker(engine, expire_on_commit=False)
+ async with engine.begin() as conn:
+ await conn.run_sync(Base.metadata.create_all)
+
+ now = datetime.now(UTC)
+ async with session_factory() as db:
+ db.add(_window_test_source(now))
+ await _seed_window_items(db, now)
+
+ repo = ContentRepo(db)
+ items = await repo.list_for_today_picks(hours=24)
+
+ assert {item.id for item in items} == {2, 3}, "3 天前发布的旧文即使刚被抓到也不应进 24h 候选窗口"
+
+ await engine.dispose()
+
+
+@pytest.mark.asyncio
+async def test_report_window_prefers_published_at():
+ """报告窗口按原文发布时间归档:晚抓到的旧文不进之后的日报。"""
+ engine = create_async_engine("sqlite+aiosqlite:///:memory:")
+ session_factory = async_sessionmaker(engine, expire_on_commit=False)
+ async with engine.begin() as conn:
+ await conn.run_sync(Base.metadata.create_all)
+
+ now = datetime.now(UTC)
+ async with session_factory() as db:
+ db.add(_window_test_source(now))
+ await _seed_window_items(db, now)
+
+ repo = ContentRepo(db)
+ # 报告窗口覆盖「刚刚」(旧文 id=1 的 crawled_at 在窗口内、published_at 在窗口外)
+ report_items = await repo.list_for_report_window(
+ window_start=now - timedelta(hours=1),
+ window_end=now + timedelta(hours=1),
+ )
+ assert 1 not in {
+ item.id for item in report_items
+ }, "旧文(published_at 在报告窗口外)不应因晚抓而出现在本期日报"
+ assert 3 in {item.id for item in report_items}
+
+ await engine.dispose()
diff --git a/docs/quality/regression-matrix.md b/docs/quality/regression-matrix.md
index aafc1b8a..b3aaa758 100644
--- a/docs/quality/regression-matrix.md
+++ b/docs/quality/regression-matrix.md
@@ -89,7 +89,7 @@
- 症状:`_cache_warmup_task` 非 CancelledError 异常会中断后续全部清理步骤;jieba 预热 await 无超时且 `to_thread` 不可取消,可挂死停机;整体停机无 deadline。
- 回归测试:`tests/test_shutdown_prewarm.py`(异常不外抛且留痕 / jieba 超时不挂死 / 正常与已取消路径)。owner:#73(Parent #6,修复 PR #77:`_shutdown_prewarm_tasks`)。
-### 管理后台批量操作静默失败 — #88 🔧 修复完成待合并(2026-09-29)
+### 管理后台批量操作静默失败 — #88 ✅ 已修复关闭(2026-09-29,PR #89 合并)
- 症状:`sources/page.tsx` `handleBatchToggle` 逐条 `sourcesApi.update` 的 `catch` 只写 `console.error`,循环结束后**无条件** `setSelectedIds(new Set())` + `fetchSources()`。批量停用 30 个信源、其中 5 个失败时,UI 表现为「全部处理完」,失败的 5 个仍按原状态继续采集并污染内容池,而操作者无从察觉。
- 环境/前提:生产数据,仅管理员账号;`PATCH /api/v1/sources/{id}` 任一请求失败(网络抖动 / 404 / 5xx)即可触发。
- 复现:信源管理页多选 ≥2 个信源 → 批量停用 → 中途阻断网络或使其中之一更新失败 → 观察 UI 与浏览器控制台对照。
@@ -99,7 +99,7 @@
- 回归测试:`npx vitest run src/app/admin/sources/_batch-utils.test.ts`(11 条)。选择态重算与结果文案已抽为纯函数 `_batch-utils.ts`,`.tsx` 只做调用;已用变异测试验证——把实现改回旧行为(失败进 console + 无条件清空)后 3 条转红,其中包含「中途新勾选项被静默丢弃」这条。
- owner:#88。
-### Webhook 日志页读数口径撒谎 — #88 🔧 修复完成待合并(2026-09-29)
+### Webhook 日志页读数口径撒谎 — #88 ✅ 已修复关闭(2026-09-29,PR #89 合并)
- 症状:`webhook-logs/page.tsx` 的 `successCount` / `failCount` 只统计**当前页** 30 行,却与全局 `total` 并排渲染成同款 Badge。读者会把「本页 2 失败」除以「共 1240 条」读成 0.16% 失败率,真相是第 5 页可能还躺着 40 条。更严重的是失败徽章以 `failCount > 0` 为条件渲染——翻到失败为 0 的分页时「失败」整枚徽章消失,**徽章的缺席本身制造错误信念**。
- 环境/前提:任何有 ≥1 页推送日志的生产数据。
- 复现:`/admin/webhook-logs` 翻页,观察任意分页顶部的 Badge 组。
@@ -110,7 +110,7 @@
- **覆盖边界(独立复核 M3 实证,勿夸大)**:修复前 `summarizeLogPage` 的 6 个公式与修复后逐字相同,只测它**证明不了**本缺陷;真正的修复是把「渲染成什么」下沉成 `buildSummaryBadges`,其断言(失败徽章 0 时仍存在且为中性色、全部徽章带「本页」/「全部」口径词、空页只留 1 枚)才真正对应缺陷本身。JSX 渲染层仍无组件测试(见下方遗留项)。
- owner:#88。
-### 管理后台面包屑缺 4 项映射 — #88 🔧 修复完成待合并(2026-09-29)
+### 管理后台面包屑缺 4 项映射 — #88 ✅ 已修复关闭(2026-09-29,PR #89 合并)
- 症状:`AdminTopBar.tsx` 的 `ADMIN_PAGE_LABELS` 只有 10 条,侧边栏 `ADMIN_NAV_ITEMS` 有 15 项。`prompts` / `scoring-dashboard` / `evidence` / `webhook-logs` 四页无映射,`findPageLabel` 回退显示兜底文案「管理」,与侧边栏自相矛盾。(另经核实:概览卡片 13 张、面包屑 10 条、侧边栏 15 项三份目录互相矛盾,说明分类从未被写下来过。)
- 环境/前提:always。
- 复现:访问上述四页之一,观察顶栏面包屑第二段。
@@ -120,6 +120,17 @@
- 回归测试:`npx vitest run src/lib/__tests__/admin-nav.test.ts`(13 条)。导航清单与面包屑映射合并到 `src/lib/admin-nav.ts` 单一事实源,映射由清单派生;`nav-checklist` 等价断言(壳内每项都有映射、label 与清单逐项一致)已在其中,漏项即红。顺带修掉一个被测试抓出的旧行为:`/admin/任何未知路径` 曾因壳根前缀匹配被判成「概览」,现已回退兜底。
- owner:#88。
+### LLM 降级内容再也不会被重新分析:local_fallback 无恢复路径 — #90 🔧 修复完成待合并(2026-10-01)
+- 症状:LLM 降级路径(熔断 `CircuitOpenError` / 内容过滤 `BadRequestError` / 预算 `LlmBudgetExceededError` 三种触发源共用)把内容写为 `ANALYZED` 终态、`AiAnalysis.summary_source='local_fallback'` 落库后,**全代码库没有任何 requeue 消费者**——重分析资格谓词只认 PENDING/僵死 ANALYZING/ERROR,手动分析端点遇已有记录早退。降级内容是确定性假分数(curation 61-64 窄带、risk 恒 28、模板推荐语),today-picks 以 40px 大字渲染且前端不读 `summary_source`,用户无法分辨;窄带假分参与 P70 计算会拖低百分位门槛、挤掉真实 65-70 分内容。
+- 环境/前提:任一降级触发源命中即触发;当前部署全走本地模型(`llm_call_logs.total_cost` 全 0),预算闸实际不触发,触发面在切换付费 API 后成为现实。
+- 复现:触发降级(如断言 `LlmBudgetExceededError` 路径或熔断 OPEN)→ 确认内容 `status=ANALYZED`、`summary_source='local_fallback'` → 检索全部重分析入口(`_analysis_candidate_condition`、`api/v1/analyses.py` 手动端点、scripts/)确认无路径捡起该内容。
+- 期望 vs 实际:期望触发原因消除后(预算窗口滚过/熔断恢复)内容可被重新 LLM 分析;实际永久停留假分形态。
+- 边界:`analysis.py` 降级路径 + `content_repo.py` 重分析资格谓词。已修(分支 `wip-llm-budget-guard`,2 笔提交):新增 `services/analysis_requeue.py` requeue 服务 + scheduler 每 15 分钟调度,预算余量 ≥30% 且降级满 60 分钟后把 `local_fallback` 内容重置回 PENDING;P70 门槛改为只由真实 LLM 评分决定(`scoring_engine.py`);`summary_source` 透传补齐 DuckDB 主路径、OLTP fallback、API payload 与周/月报 digest;today-picks 卡片显示「本地速览 · 待 AI 复核」标记;预算闸豁免 daily_report/weekly_digest/monthly_digest。
+- 严重度:P1(预算日窗最坏锁 ~24h × 实测 ~250 条/时摄入 ≈ 数千条永久假分;当前零账单部署不触发故非 P0);**复现性**:always(降级发生即永久停留假分)。
+- 回归测试:`uv run python -m pytest tests/test_analysis_requeue.py -q`(6 条:资格与字段重置 / limit 最旧优先 / 服务层余量门控 / 余量阈值语义 / 豁免场景绕行 / ≥2 条 fallback 不再回收);`tests/test_llm_budget_guard.py`(9 条);`tests/test_scoring_engine.py -k "p70 or all_fallback"`(2 条)。变异验证(2026-10-01 独立复核):删掉 P70 排除逻辑 → `test_p70_threshold_excludes_local_fallback_fake_scores` 转红;删掉全降级批次的 `or [...]` 回退分支 → `test_all_fallback_batch_falls_back_to_full_scores` 转红(**该分支非防御性冗余**:真实分为空时 `_compute_percentile_threshold` 返回全局默认 55,而降级内容 final_score 普遍低于 55,回退缺失会导致整批选不出内容)。全量后端套件仅 2 条 FAILED,均为 `test_source_url_safety.py` 依赖真实网络的存量红,在 main 上同样 FAILED,非本次引入。
+- owner:#90。
+- **存量定性(勿误读)**:非预算闸(`feat(llm): 增加三窗预算熔断闸`)引入——熔断降级早就同样留下假分内容;预算闸放大触发面,圆桌评审三席独立发现。发现方式:预算数值圆桌(2026-09-30,Memory 存 `.vidt/roundtable/budget-fuse-values/`,会话态)。
+
## 三、关键流程基线(9 项)
状态标记:✅ = 2026-09-27 在 main @ 7203847 新鲜复跑通过;📋 = 现有套件覆盖、未逐项复跑(跑全量即覆盖);🔧 = 修复完成待合并(PR 已开、CI 绿、独立复核非 pass 尚未全部处置)。
diff --git a/frontend/src/app/today-picks/_components.tsx b/frontend/src/app/today-picks/_components.tsx
index c2d44d5e..5a562471 100644
--- a/frontend/src/app/today-picks/_components.tsx
+++ b/frontend/src/app/today-picks/_components.tsx
@@ -630,6 +630,14 @@ export function PickCard({
/
{timeAgo(item.published_at || item.crawled_at)}
+ {analysis?.summary_source === 'local_fallback' && (
+
+ 本地速览 · 待 AI 复核
+
+ )}
{item.category && {item.category}}
{item.content_type && {item.content_type}}
{tags.slice(0, 3).map((tag) => (
diff --git a/frontend/src/types/index.ts b/frontend/src/types/index.ts
index a27eb3c8..7cc52682 100644
--- a/frontend/src/types/index.ts
+++ b/frontend/src/types/index.ts
@@ -340,6 +340,8 @@ export interface ContentAnalysis {
creator_score: number;
viral_score: number;
risk_score: number;
+ // 分析来源标记:'local_fallback' = 降级本地速览(#90,前端据此显示待复核标记)
+ summary_source?: string | null;
platform_fit?: Record | null;
recommended_reason?: string | null;
summary?: string | null;