Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions .gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -73,3 +73,6 @@ htmlcov/
.brooks-lint-history.json
.vidt/
.node-version

# 独立 worktree 工作区(issue 切片隔离)
.worktrees/
19 changes: 19 additions & 0 deletions backend/.env.example
Original file line number Diff line number Diff line change
Expand Up @@ -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 渲染对应按钮。
Expand Down
10 changes: 10 additions & 0 deletions backend/app/core/config.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
81 changes: 78 additions & 3 deletions backend/app/repositories/content_repo.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -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:
Expand Down Expand Up @@ -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())
Expand Down
24 changes: 24 additions & 0 deletions backend/app/scheduler.py
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand Down Expand Up @@ -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(
Expand Down
2 changes: 1 addition & 1 deletion backend/app/services/_duckdb_picks_mixin.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
12 changes: 8 additions & 4 deletions backend/app/services/_duckdb_reports_mixin.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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
]
Expand Down
14 changes: 9 additions & 5 deletions backend/app/services/analysis.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
47 changes: 47 additions & 0 deletions backend/app/services/analysis_requeue.py
Original file line number Diff line number Diff line change
@@ -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
3 changes: 3 additions & 0 deletions backend/app/services/creation.py
Original file line number Diff line number Diff line change
Expand Up @@ -26,13 +26,16 @@

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

if isinstance(exc, LlmCapacityUnavailableError):
return "创作方案暂时排队等待可用模型渠道,请稍后重试。"
if isinstance(exc, CircuitOpenError):
return "创作方案服务暂时不可用,系统正在自动恢复,请稍后重试。"
if isinstance(exc, LlmBudgetExceededError):
return "模型调用预算已达今日上限,创作方案稍后自动恢复,请明日再试。"
return "创作方案生成失败,请稍后重试。"


Expand Down
1 change: 1 addition & 0 deletions backend/app/services/digest_context.py
Original file line number Diff line number Diff line change
Expand Up @@ -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),
Expand Down
5 changes: 4 additions & 1 deletion backend/app/services/llm/_call_engine.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)

Expand Down
Loading
Loading