Author SHA1 Message Date
WorkBuddy d9aeff2114 debt(D7): 限流/KeyRing 单进程限制说明 + 生产启动一次性告警
- app.py lifespan:APP_ENV=production 时 logger.warning 提醒
  进程内 rate-limit/KeyRing 仅单进程有效(不 sys.exit、不引入 Redis)
- key_ring.py docstring 补充多 worker 不共享的后果说明
- deps.py 注释已具备(单进程有效,生产前置 Nginx),无需改动
2026-09-21 19:51:41 +08:00
WorkBuddy 17e99a4f9a debt(D5): agent_outputs 类型标注对齐真实 JSON 形状(list[dict]|dict|None)
- Mapped[list[dict] | dict | None]:multi 模式存专家报告列表,
  历史数据/兼容路径可能存 dict;加注释说明形状来源
- 仅类型标注修正,列类型(JSONB)与数据、迁移均不变
2026-09-21 19:51:38 +08:00
WorkBuddy 4514ef4e92 debt(D4): bzzoiro Team/League/Match 查找经 Repository 层,纠正空壳表述
- repositories.py 新增 MatchRepository.find_by_league_and_date_range /
  find_finished_with_stats(含 stats 预加载),TeamRepository.get_or_create
  支持 name_zh
- bzzoiro events/standings/stats 三条管线的联赛、球队、比赛查找与批量预载
  全部改走 Repository;事务提交仍由调用方 UnitOfWork 控制
- Standing/RawEvent/Lineage 管线内私有读写保留在本模块(不强行 Repository 化)
- 模块 docstring 纠正为准确表述,移除「已全面 Repository 化」的误导性声明
2026-09-21 19:48:56 +08:00
WorkBuddy 18c89111d8 debt(D6): Admin 导航单源 —— NAV_ITEMS 派生侧栏/命令面板/面包屑
- 新增 admin/nav.ts: 唯一配置源 NAV_ITEMS(to/label/group/icon/hideFromSidebar)
- COMMAND_PALETTE_PAGES / ROUTE_LABELS / NAV_SECTIONS 全部由其派生
- 删除 AdminLayout 中三份平行清单(已出现漂移: data-pipeline 不在侧栏、
  「系统」vs「系统设置」命名不一致、monitoring/eval 顺序两处相反)
- 已知微差: 命令面板「评估与监控」组内顺序统一为侧栏口径(eval 在前)
- 行为不变: 路由/标签/分组标题均保持现状;data-pipeline 保持仅面板可达
2026-09-21 19:34:31 +08:00
WorkBuddy 52d67863d6 debt(D3): 拆分 Matches.tsx(~1400 行)为 hooks + 组件 + 共享层
- pages/matches/types.ts: Match/Prediction/AgentReport + 常量(LEAGUES/STATUS_META/AGENT_LABELS 等)
- pages/matches/ui.tsx: Spinner/SkeletonRows/Switch/日期分组工具(Switch 自组件体内提升到模块级)
- pages/matches/hooks/useMatchesList.ts: 列表筛选/游标分页/进行中比赛,竞态防护语义不变
- pages/matches/hooks/useMatchPredict.ts: 预测状态机(连点/竞态/中止/超时)+ 可读错误文案
- pages/matches/components/MatchPredictPanel.tsx: 预测弹窗全套(过程可视化/结果/专家意见)
- pages/matches/components/MatchDetailSection.tsx: MatchRow(原内联行 JSX 抽出)+ 详情面板
- Matches.tsx 组装页 279 行(< 300);error state 共享行为与拆分前一致
- 验证: tsc + vite build 通过,渲染逻辑原样搬迁
2026-09-21 19:31:42 +08:00
WorkBuddy 5c7fdce0a3 debt(D2): 统一预测结果类型为 PredictResult,路由去 dict 分支
- PredictResult 扩展可选字段 mode/agent_outputs/agent_weights/prompt_tokens/completion_tokens
- MultiPredictResult 变为 PredictResult 别名(保留 R4 守卫标记与 re-export)
- baseline 改返回 PredictResult(修复 backtest 对 baseline AttributeError 的潜伏 bug)
- 预测路由单一属性映射,删除全部 isinstance(result, dict) 分支
- _persist_baseline 属性化,baseline upsert 语义不变(prompt_version/token/latency 同前)
- TDD: 5 新测试 + test_baseline.py 属性化;全量 257 passed
2026-09-21 19:23:50 +08:00
WorkBuddy 3a9f3f5a0e debt(D1): events 成功路径补写 Bronze 层(RawEvent + DataLineage)
- 新增 _events_record_id: 上游 id 缺失时用 (league:home:away:date) 合成稳定幂等键
- 新增 _write_events_bronze: best-effort 写 RawEvent(幂等) + Lineage(matches/events_ingest)
- ingest 插入与变更更新后触发;同批 seen 集合防重复;基础设施失败只 warning
- TDD: 6 测试(插入/合成键/幂等跳过/变更更新/无变化不写/失败不拖垮),双变异验证通过
2026-09-21 19:23:37 +08:00
shangfangjian f6c0145c32 Merge pull request 'fix: 死信表真正接线 + 前端 HTTP 收敛与状态修正' (#9) from fix-deadletter-frontend-cleanup into main
Reviewed-on: #9
2026-09-21 18:39:08 +08:00
WorkBuddy b6a251407d fix(pipeline): 死信表真正接线 + 前端 HTTP 收敛与状态修正
- bzzoiro events/standings/stats 抓取失败写入 IngestFailure 死信(尽力而为,
  写入失败不影响主流程);顺带修复 standings 失败路径 league_r 缺 errors 键
  的 KeyError —— 该路径此前从未跑通,一旦失败会顶掉原始异常
- admin/api.ts 收敛为 lib/http.ts 薄门面,消除第二套 HTTP 实现;
  UNAUTHORIZED_EVENT 定义移至共享层,断开 lib→admin 反向依赖
- STATUS_META 死键 live 改为 in_play(对齐 normalize.py 口径),
  补 paused/postponed/cancelled/suspended;移除无人消费的
  DashboardStats.total_matches(items.length 近似,上限 100)
- 新增 tests/test_ingest_deadletter.py(6 例,变异验证判别力)
2026-09-21 18:36:43 +08:00
shangfangjian b5f7e682f0 Merge pull request 'fix: 批量修复一些问题
Reviewed-on: #8
2026-09-21 18:04:54 +08:00
26 changed files with 2464 additions and 1487 deletions
+5 -53
View File
@@ -9,6 +9,7 @@ import { useState, useEffect, useCallback } from 'react'
import { NavLink, Outlet, useLocation } from 'react-router-dom' import { NavLink, Outlet, useLocation } from 'react-router-dom'
import { fetchAuthState, logout, UNAUTHORIZED_EVENT } from './api' import { fetchAuthState, logout, UNAUTHORIZED_EVENT } from './api'
import { fetchHealth } from './dal' import { fetchHealth } from './dal'
import { COMMAND_PALETTE_PAGES, ROUTE_LABELS, NAV_SECTIONS } from './nav'
import Login from './Login' import Login from './Login'
import { useCommandPalette, CommandPalette } from './useCommandPalette' import { useCommandPalette, CommandPalette } from './useCommandPalette'
@@ -88,58 +89,9 @@ function Icon({ name }: { name: string }) {
} }
} }
const NAV_PAGES: Array<{ to: string; label: string; group: string }> = [ // D6: 导航三视图(侧栏/命令面板/面包屑)统一由 admin/nav.ts 的 NAV_ITEMS
{ to: '/admin', label: '仪表盘', group: '概览' }, // 单源派生 —— 此处的 NAV_PAGES/ROUTE_LABELS/NAV_SECTIONS 平行清单已删除,
{ to: '/admin/collection', label: '数据采集', group: '数据流水线' }, // 新增页面/改标签只改 nav.ts 一处。
{ to: '/admin/data-completeness', label: '数据完整性', group: '数据流水线' },
{ to: '/admin/data-pipeline', label: '数据管线', group: '数据流水线' },
{ to: '/admin/predictions', label: '预测历史', group: '数据流水线' },
{ to: '/admin/backtest', label: '回测', group: '数据流水线' },
{ to: '/admin/monitoring', label: '监控', group: '评估与监控' },
{ to: '/admin/eval', label: '评估', group: '评估与监控' },
{ to: '/admin/settings', label: '设置', group: '系统' },
{ to: '/admin/logs', label: '日志', group: '系统' },
]
// 路由 → 面包屑标签
const ROUTE_LABELS: Record<string, string> = {
'/admin': '仪表盘',
'/admin/collection': '数据采集',
'/admin/data-completeness': '数据完整性',
'/admin/data-pipeline': '数据管线',
'/admin/predictions': '预测历史',
'/admin/backtest': '回测',
'/admin/monitoring': '监控',
'/admin/settings': '设置',
'/admin/logs': '日志',
'/admin/eval': '评估',
}
const NAV_SECTIONS: { title: string; items: Array<{ to: string; label: string; icon: string }> }[] = [
{
title: '数据流水线',
items: [
{ to: '/admin/collection', label: '数据采集', icon: 'collection' },
{ to: '/admin/data-completeness', label: '数据完整性', icon: 'chart' },
{ to: '/admin/predictions', label: '预测历史', icon: 'logs' },
{ to: '/admin/backtest', label: '回测', icon: 'repeat' },
],
},
{
title: '评估与监控',
items: [
{ to: '/admin/eval', label: '评估', icon: 'eval' },
{ to: '/admin/monitoring', label: '监控', icon: 'monitor' },
],
},
{
title: '系统设置',
items: [
{ to: '/admin/settings', label: '设置', icon: 'settings' },
{ to: '/admin/logs', label: '日志', icon: 'logs' },
],
},
]
/** 报眉日期行,与前台同款式 */ /** 报眉日期行,与前台同款式 */
function dateLine(): string { function dateLine(): string {
@@ -156,7 +108,7 @@ export default function AdminLayout() {
const [healthOk, setHealthOk] = useState<boolean | null>(null) const [healthOk, setHealthOk] = useState<boolean | null>(null)
const [authed, setAuthed] = useState<boolean | null>(null) const [authed, setAuthed] = useState<boolean | null>(null)
const location = useLocation() const location = useLocation()
const palette = useCommandPalette(NAV_PAGES) const palette = useCommandPalette(COMMAND_PALETTE_PAGES)
// 登录门禁:挂载时探测会话,收到 401 事件(会话过期)自动切回登录页 // 登录门禁:挂载时探测会话,收到 401 事件(会话过期)自动切回登录页
useEffect(() => { useEffect(() => {
+20 -87
View File
@@ -1,102 +1,35 @@
/** /**
* Admin 后台管理系统 - 统一 API 客户端 * Admin 后台管理系统 - 统一 API 客户端(门面)
* *
* 鉴权:通过 POST /api/v1/auth/login 用密码换取 HttpOnly Cookie 会话, * 鉴权:通过 POST /api/v1/auth/login 用密码换取 HttpOnly Cookie 会话,
* 同源请求自动携带 Cookie,无需手动管理密钥。 * 同源请求自动携带 Cookie,无需手动管理密钥。
* 收到 401 时广播 `profeto:unauthorized` 事件,由 AdminLayout 切回登录页。 * 收到 401 时广播 `profeto:unauthorized` 事件,由 AdminLayout 切回登录页。
*
* 实现已收敛到共享层 lib/http.ts(超时/错误解析/401 广播只此一份),
* 本文件仅保留 Admin 侧的门面签名与认证接口,供既有页面按原路径导入。
*/ */
import { http, ApiError, UNAUTHORIZED_EVENT } from '../lib/http'
/** Admin 侧兼容导出:错误类型与会话失效事件名的规范来源在 lib/http */
export { ApiError, UNAUTHORIZED_EVENT }
const API_BASE = '/api/v1' const API_BASE = '/api/v1'
const TIMEOUT_MS = 30_000
/** 会话失效事件名,AdminLayout 监听后弹出登录页 */ /** Admin 请求可覆盖项(与 lib/http RequestOptions 对齐的子集) */
export const UNAUTHORIZED_EVENT = 'profeto:unauthorized' type ApiOpts = {
timeoutMs?: number
export class ApiError extends Error { /** 改密接口的 401 表示「当前密码错误」,非会话过期,置 true 跳过登出广播 */
constructor( skipAuthHandling?: boolean
message: string,
public status: number,
public data?: unknown,
) {
super(message)
this.name = 'ApiError'
}
}
async function request<T>(
path: string,
options: RequestInit & { timeoutMs?: number; skipAuthHandling?: boolean } = {},
): Promise<T> {
// 修复: 正确拼接 API_BASE
const url = path.startsWith('http')
? path
: path.startsWith('/')
? path // 已经是绝对路径(如 /health)
: `${API_BASE}${path}`
const { timeoutMs = TIMEOUT_MS, skipAuthHandling, ...fetchOptions } = options
const controller = new AbortController()
const timer = setTimeout(() => controller.abort(), timeoutMs)
try {
const res = await fetch(url, {
...fetchOptions,
signal: controller.signal,
headers: {
'Content-Type': 'application/json',
...fetchOptions.headers,
},
})
if (!res.ok) {
// 先读 text 再尝试 JSON 解析:Response body 流只能读一次,
// 若先调 res.json() 失败(如返回 HTML 错误页),再调 res.text() 会抛 "body stream already read"。
const rawText = await res.text()
let detail: unknown = rawText
try {
detail = JSON.parse(rawText)
} catch {
// 非 JSON(如 HTML 错误页),保留原始文本
}
let message =
detail && typeof detail === 'object' && detail !== null && 'detail' in detail
? String((detail as { detail: unknown }).detail)
: `HTTP ${res.status}: ${res.statusText}`
// 401 仅在没有显式跳过时广播未登录事件(改密接口的 401 表示当前密码错误,非会话过期)
if (res.status === 401 && !skipAuthHandling) {
message += '\n登录已过期,请重新登录。'
window.dispatchEvent(new CustomEvent(UNAUTHORIZED_EVENT))
}
throw new ApiError(message, res.status, detail)
}
// 修复: 正确判断 204 No Content
if (res.status === 204) {
return undefined as T
}
return res.json() as Promise<T>
} catch (err) {
if (err instanceof ApiError) throw err
if (err instanceof DOMException && err.name === 'AbortError') {
throw new ApiError('请求超时,请稍后重试', 0)
}
throw new ApiError(
err instanceof Error ? err.message : '网络错误,请检查连接',
0,
)
} finally {
clearTimeout(timer)
}
} }
export const api = { export const api = {
get: <T>(path: string) => request<T>(path), get: <T>(path: string) => http.get<T>(path),
post: <T>(path: string, body?: unknown, opts?: { timeoutMs?: number; skipAuthHandling?: boolean }) => post: <T>(path: string, body?: unknown, opts?: ApiOpts) =>
request<T>(path, { method: 'POST', body: body ? JSON.stringify(body) : undefined, ...opts }), http.post<T>(path, body, opts),
put: <T>(path: string, body?: unknown, opts?: { timeoutMs?: number; skipAuthHandling?: boolean }) => put: <T>(path: string, body?: unknown, opts?: ApiOpts) =>
request<T>(path, { method: 'PUT', body: body ? JSON.stringify(body) : undefined, ...opts }), http.put<T>(path, body, opts),
delete: <T>(path: string) => request<T>(path, { method: 'DELETE' }), delete: <T>(path: string) => http.delete<T>(path),
} }
// ── 认证 ──────────────────────────────────────────────────────── // ── 认证 ────────────────────────────────────────────────────────
+5 -7
View File
@@ -32,23 +32,21 @@ import type {
* 从多个端点聚合仪表盘数据。 * 从多个端点聚合仪表盘数据。
* 后端暂无专用仪表盘端点,这里组合 health + 各列表端点。 * 后端暂无专用仪表盘端点,这里组合 health + 各列表端点。
* *
* P2-8 修复: 使用 items.length 替代不存在的 total 字段, * 注: 比赛真实总量请用 fetchAdminStats()(GET /admin/stats,
* 并扩大 limit 以获得更有参考价值的数量。 * 后端 COUNT(*) 精确计数)。此处曾用 matches items.length 近似,
* 已随仪表盘切换真实计数而移除,防止误用 100 上限的假总量。
*/ */
export async function fetchDashboard(): Promise<DashboardStats> { export async function fetchDashboard(): Promise<DashboardStats> {
// 并行获取各端点数据 // 并行获取各端点数据
// matches 返回 {items, next_cursor, has_more}, predictions 返回数组 // predictions 返回数组
const [leagues, matches, predictions, health] = await Promise.allSettled([ const [leagues, predictions, health] = await Promise.allSettled([
api.get<League[]>(`${API_BASE}/leagues`), api.get<League[]>(`${API_BASE}/leagues`),
api.get<{ items: Match[]; has_more: boolean }>(`${API_BASE}/matches?limit=100`),
api.get<Prediction[]>(`${API_BASE}/predictions?limit=100`), api.get<Prediction[]>(`${API_BASE}/predictions?limit=100`),
api.get<{ status: string }>('/health'), api.get<{ status: string }>('/health'),
]) ])
return { return {
leagues: leagues.status === 'fulfilled' ? leagues.value : [], leagues: leagues.status === 'fulfilled' ? leagues.value : [],
// P2-8: matches 无 total 字段,用 items.length 近似(上限 100)
total_matches: matches.status === 'fulfilled' ? (matches.value as any)?.items?.length ?? 0 : 0,
// predictions 直接返回数组 // predictions 直接返回数组
total_predictions: predictions.status === 'fulfilled' ? (predictions.value as any)?.length ?? 0 : 0, total_predictions: predictions.status === 'fulfilled' ? (predictions.value as any)?.length ?? 0 : 0,
health: health.status === 'fulfilled' ? (health.value as any).status : 'unknown', health: health.status === 'fulfilled' ? (health.value as any).status : 'unknown',
+67
View File
@@ -0,0 +1,67 @@
/**
* Admin 导航单源(D6)。
*
* 背景: 侧栏(NAV_SECTIONS)、命令面板(NAV_PAGES)、面包屑(ROUTE_LABELS)
* 此前各维护一份平行清单,已出现漂移 —— data-pipeline 不在侧栏、
* 「系统」vs「系统设置」命名不一致、monitoring/eval 两处顺序相反。
*
* 现在只允许维护 NAV_ITEMS 一份,其余视图一律由此派生;
* 禁止再新建平行导航清单(修改入口/新增页面只改这里)。
*/
export interface NavItem {
to: string
label: string
/** 命令面板中的分组名(展示原样) */
group: string
/** 侧栏图标名(见 AdminLayout 的 Icon) */
icon: string
/** 仅命令面板/面包屑可达,不进侧栏(深链页) */
hideFromSidebar?: boolean
}
/** 唯一的导航配置源。顺序 = 侧栏渲染顺序(命令面板分组内顺序与之相同)。 */
export const NAV_ITEMS: NavItem[] = [
{ to: '/admin', label: '仪表盘', group: '概览', icon: 'chart' },
{ to: '/admin/collection', label: '数据采集', group: '数据流水线', icon: 'collection' },
{ to: '/admin/data-completeness', label: '数据完整性', group: '数据流水线', icon: 'chart' },
{ to: '/admin/data-pipeline', label: '数据管线', group: '数据流水线', icon: 'chart', hideFromSidebar: true },
{ to: '/admin/predictions', label: '预测历史', group: '数据流水线', icon: 'logs' },
{ to: '/admin/backtest', label: '回测', group: '数据流水线', icon: 'repeat' },
{ to: '/admin/eval', label: '评估', group: '评估与监控', icon: 'eval' },
{ to: '/admin/monitoring', label: '监控', group: '评估与监控', icon: 'monitor' },
{ to: '/admin/settings', label: '设置', group: '系统', icon: 'settings' },
{ to: '/admin/logs', label: '日志', group: '系统', icon: 'logs' },
]
/** 命令面板条目(原 NAV_PAGES 的唯一来源) */
export const COMMAND_PALETTE_PAGES = NAV_ITEMS.map(({ to, label, group }) => ({ to, label, group }))
/** 路由 → 面包屑标签(原 ROUTE_LABELS 的唯一来源) */
export const ROUTE_LABELS: Record<string, string> = Object.fromEntries(
NAV_ITEMS.map(i => [i.to, i.label]),
)
/** 侧栏分组标题的历史显示名(仅侧栏使用;与面板分组名不同时在此映射) */
const SIDEBAR_GROUP_TITLES: Record<string, string> = {
: '系统设置',
}
/**
* 侧栏分组(原 NAV_SECTIONS 的唯一来源)。
* 仪表盘(group=概览)在 AdminLayout 中独立渲染于顶部,不进分组循环。
*/
export const NAV_SECTIONS = (() => {
const sidebarItems = NAV_ITEMS.filter(i => !i.hideFromSidebar && i.group !== '概览')
const titles: string[] = []
for (const i of sidebarItems) {
const title = SIDEBAR_GROUP_TITLES[i.group] ?? i.group
if (!titles.includes(title)) titles.push(title)
}
return titles.map(title => ({
title,
items: sidebarItems
.filter(i => (SIDEBAR_GROUP_TITLES[i.group] ?? i.group) === title)
.map(({ to, label, icon }) => ({ to, label, icon })),
}))
})()
+2 -1
View File
@@ -17,7 +17,8 @@ export interface HealthStatus {
export interface DashboardStats { export interface DashboardStats {
leagues: League[] leagues: League[]
total_matches: number // 比赛总量已移除: 用 fetchAdminStats()(/admin/stats)的精确 COUNT,
// 不要再用列表 items.length 近似(上限 100 会严重失真)
total_predictions: number total_predictions: number
health: string health: string
db_tables: { name: string; row_count: number; size_mb: number; last_updated: string | null }[] db_tables: { name: string; row_count: number; size_mb: number; last_updated: string | null }[]
+4 -1
View File
@@ -7,9 +7,12 @@
* - 401 自动广播(Admin 场景) * - 401 自动广播(Admin 场景)
* - JSON/HTML 容错解析 * - JSON/HTML 容错解析
* - 请求竞态防护(可选 signal) * - 请求竞态防护(可选 signal)
*
* 注: Admin 侧的 admin/api.ts 是本模块的薄门面,不再有第二套实现。
*/ */
import { UNAUTHORIZED_EVENT } from '../admin/api' /** 会话失效事件名,AdminLayout 监听后弹出登录页 */
export const UNAUTHORIZED_EVENT = 'profeto:unauthorized'
const API_BASE = '/api/v1' const API_BASE = '/api/v1'
const DEFAULT_TIMEOUT = 30_000 const DEFAULT_TIMEOUT = 30_000
File diff suppressed because it is too large Load Diff
@@ -0,0 +1,369 @@
/**
* MatchDetailSection: 赛程行(MatchRow)+ 展开详情面板(统计/近况/H2H/历史预测)。
*
* D3: 从 Matches.tsx 拆出,渲染逻辑原样搬迁。MatchRow 原先是主组件内
* group.map 的内联 JSX,现接收回调(onToggle/onPredict)保持行为一致;
* 展开懒加载的 state 仍由页面持有。
*/
import TeamSideTag from '../../../components/TeamSideTag'
import { fetchMatchDetail, fetchMatchContext } from '../../../admin/dal'
import type { MatchDetailOut, MatchContextOut, MatchRecentPrediction, TeamRecentMatch } from '../../../admin/types'
import type { MatchStatsDetail } from '../../../admin/types'
import { STATUS_META } from '../types'
import type { Match } from '../types'
import { Spinner } from '../ui'
/** 比赛详情面板:双方近况/H2H + 历史预测列表(只读) */
function MatchDetailPanel({
match, detail, ctx, loading,
}: {
match: Match
detail: MatchDetailOut | undefined
ctx: MatchContextOut | undefined
loading: boolean
}) {
const homeName = match.home_team_zh || match.home_team
const awayName = match.away_team_zh || match.away_team
const finished = match.match_status === 'finished'
return (
<div className="border-b border-ink-200 bg-paper-100/50 px-3 py-4">
{loading && (
<div className="flex items-center gap-2 text-xs text-ink-500"><Spinner /> </div>
)}
{!loading && !detail && !ctx && (
<p className="py-4 text-center text-xs text-ink-400"></p>
)}
{!loading && (detail || ctx) && (
<div className="space-y-5">
{/* 比分区(终场/当前比分 + 状态 + 预测按钮) */}
<div className="flex flex-wrap items-center justify-between gap-3">
<div className="text-center">
<p className="font-serif text-3xl font-bold tabular-nums leading-none text-ink-900">
{match.home_goals ?? '-'}{' '}<span className="text-ink-300">:</span>{' '}{match.away_goals ?? '-'}
</p>
<p className="mt-1 text-2xs text-ink-500">
{match.match_stage || ''} {match.match_status === 'finished' ? '· 已完赛' : match.match_status === 'scheduled' ? '· 未开赛' : `· ${match.match_status}`}
</p>
{match.home_xg != null && match.away_xg != null && (
<p className="text-2xs tabular-nums text-ink-400">xG {match.home_xg.toFixed(1)}{match.away_xg.toFixed(1)}</p>
)}
</div>
{!finished && (
<span className="text-2xs text-ink-500">
</span>
)}
</div>
{/* 比赛详细统计(bzzoiro /events/{id}/stats/) */}
{detail?.stats && (
<MatchStatsPanel stats={detail.stats} homeName={homeName} awayName={awayName} />
)}
{/* 双方近况 + H2H */}
{(ctx?.home_recent?.length || ctx?.away_recent?.length || ctx?.h2h?.length) ? (
<div className="grid gap-4 sm:grid-cols-3">
<RecentBlock title={`${homeName} 近况`} rows={ctx?.home_recent} side="home" />
<RecentBlock title={`${awayName} 近况`} rows={ctx?.away_recent} side="away" />
<RecentBlock title="历史交锋(H2H)" rows={ctx?.h2h} side="h2h" />
</div>
) : (
!loading && <p className="text-2xs text-ink-400"></p>
)}
{/* 历史预测列表 */}
<div>
<h4 className="section-head mb-2">({detail?.recent_predictions?.length ?? 0})</h4>
{detail?.recent_predictions?.length ? (
<div className="space-y-2">
{detail.recent_predictions.map(p => (
<PredictionHistoryRow key={p.id} p={p} />
))}
</div>
) : (
<p className="py-3 text-center text-2xs text-ink-400"></p>
)}
</div>
</div>
)}
</div>
)
}
/** 赛程表一行:可点击展开详情;展开时懒加载详情(只读,不触发 LLM) */
export function MatchRow({
m,
busy,
expanded,
detail,
ctx,
detailLoading,
onToggle,
onPredict,
}: {
m: Match
busy: boolean
expanded: boolean
detail: MatchDetailOut | undefined
ctx: MatchContextOut | undefined
detailLoading: boolean
onToggle: () => void
onPredict: (m: Match) => void
}) {
const st = STATUS_META[m.match_status] ?? { label: m.match_status, cls: 'text-ink-400' }
const homeName = m.home_team_zh || m.home_team
const awayName = m.away_team_zh || m.away_team
const finished = m.match_status === 'finished'
return (
<div>
{/* 行:可点击展开 */}
<div
className={`border-b border-ink-200 px-4 py-5 transition-colors hover:bg-paper-100/70 cursor-pointer sm:px-1 sm:py-4 ${
expanded ? 'bg-paper-100/60' : ''
}`}
onClick={onToggle}
role="button"
tabIndex={0}
onKeyDown={e => { if (e.key === 'Enter') onToggle() }}
aria-expanded={expanded}
>
{/* 桌面 grid: 日期 | 主队 | 比分 | 客队 | 状态 | 按钮 */}
<div className="flex flex-col gap-3 sm:grid sm:grid-cols-[96px_minmax(0,1fr)_72px_minmax(0,1fr)_64px_88px] sm:items-center sm:gap-x-4 sm:gap-y-0">
{/* 日期 + 状态:小屏同行;桌面 date 单独一列 */}
<div className="flex items-center justify-between text-xs sm:contents">
<span className="tabular-nums text-ink-500 sm:text-xs">{fmtTime(m.match_date)}</span>
<span className={`sm:hidden ${st.cls}`}>{st.label}</span>
</div>
{/* 主队 + 比分 + 客队:移动端 grid 三列(严格居中),桌面端 grid 分列 */}
<div className="grid grid-cols-[1fr_auto_1fr] items-center gap-4 sm:contents">
{/* 主队(右对齐) */}
<span className="flex min-w-0 items-center justify-end gap-2">
<TeamSideTag side="home" />
<span className="truncate text-sm font-medium text-ink-900">{homeName}</span>
</span>
{/* 比分 / VS(严格居中) */}
<span className="flex flex-col items-center justify-center">
{m.home_goals !== null && m.away_goals !== null ? (
<span className="font-serif text-xl font-bold tabular-nums leading-none text-ink-900 sm:text-xl">
{m.home_goals}<span className="mx-1 font-normal text-ink-300">:</span>{m.away_goals}
</span>
) : (
<span className="text-sm tracking-[0.2em] text-ink-500">VS</span>
)}
{m.home_xg !== null && m.away_xg !== null && (
<span className="mt-0.5 text-2xs tabular-nums text-ink-400">
xG {m.home_xg.toFixed(1)}{m.away_xg.toFixed(1)}
</span>
)}
</span>
{/* 客队(左对齐) */}
<span className="flex min-w-0 items-center gap-2">
<TeamSideTag side="away" />
<span className="truncate text-sm font-medium text-ink-900">{awayName}</span>
</span>
</div>
{/* 状态标签:小屏隐藏(已有);桌面用徽标样式 */}
<span className="hidden text-right sm:block">
<span className={`inline-block border px-1.5 py-0.5 text-2xs leading-tight ${st.cls} ${
m.match_status === 'finished'
? 'border-ink-200 text-ink-500'
: m.match_status === 'scheduled'
? 'border-ink-300 text-ink-600'
: 'border-press/30 text-press'
}`}>
{st.label}
</span>
</span>
{/* 预测按钮(统一,响应式尺寸) */}
{!finished && (
<div className="flex justify-end" onClick={e => e.stopPropagation()}>
<button
onClick={() => onPredict(m)}
disabled={busy}
className={`btn ${busy ? '' : 'btn-solid'} w-full min-h-[44px] sm:w-[84px] sm:min-h-0 sm:btn-sm`}
title="以多专家模式预测这场"
>
{busy ? (<><Spinner /> </>) : '预测'}
</button>
</div>
)}
</div>
</div>{/* 关闭可点击行 */}
{/* 展开详情面板 */}
{expanded && (
<MatchDetailPanel
match={m} detail={detail} ctx={ctx}
loading={detailLoading}
/>
)}
</div>
)
}
/** 行内时间展示:只显示 HH:mm(日期由分组头承担) */
function fmtTime(s: string): string {
const d = new Date(s)
return d.toLocaleTimeString('zh-CN', { hour: '2-digit', minute: '2-digit', hour12: false })
}
/** 比赛详细统计面板(bzzoiro /events/{id}/stats/) */
function MatchStatsPanel({
stats, homeName, awayName,
}: { stats: MatchStatsDetail; homeName: string; awayName: string }) {
const rows: Array<{ label: string; home: number | null; away: number | null; highlight?: 'high' | 'low' }> = [
{ label: '预期进球(xG)', home: stats.home_xg, away: stats.away_xg },
{ label: '射门', home: stats.home_shots, away: stats.away_shots },
{ label: '射正', home: stats.home_shots_on_target, away: stats.away_shots_on_target },
{ label: '角球', home: stats.home_corners, away: stats.away_corners },
{ label: '犯规', home: stats.home_fouls, away: stats.away_fouls },
{ label: '绝佳机会', home: stats.home_big_chances, away: stats.away_big_chances },
{ label: '黄牌', home: stats.home_yellow_cards, away: stats.away_yellow_cards },
{ label: '红牌', home: stats.home_red_cards, away: stats.away_red_cards },
]
const hasAny = rows.some(r => r.home != null || r.away != null)
if (!hasAny) return null
// 控球率用横条展示
const possHome = stats.home_possession
const possAway = possHome != null ? Math.max(0, 100 - possHome) : null
return (
<div>
<h4 className="section-head mb-2"></h4>
{/* 控球率横条 */}
{possHome != null && possAway != null && (
<div className="mb-3">
<div className="mb-1 flex justify-between text-2xs text-ink-500">
<span>{possHome.toFixed(0)}%</span>
<span className="text-ink-400"></span>
<span>{possAway.toFixed(0)}%</span>
</div>
<div className="flex h-1.5 overflow-hidden rounded-full bg-ink-200">
<div className="bg-ink-700 transition-[width] duration-500" style={{ width: `${possHome}%` }} />
<div className="bg-ink-300 transition-[width] duration-500" style={{ width: `${possAway}%` }} />
</div>
</div>
)}
{/* 主客对比表 */}
<div className="overflow-x-auto">
<table className="w-full text-xs">
<thead>
<tr className="border-b border-ink-200 text-ink-400">
<th className="py-1.5 text-left font-medium">{homeName}</th>
<th className="py-1.5 text-center font-medium text-ink-500"></th>
<th className="py-1.5 text-right font-medium">{awayName}</th>
</tr>
</thead>
<tbody>
{rows.filter(r => r.home != null || r.away != null).map(r => {
const h = r.home ?? 0
const a = r.away ?? 0
const winner = h > a ? 'home' : h < a ? 'away' : 'tie'
return (
<tr key={r.label} className="border-b border-ink-100">
<td className={`py-1.5 text-right tabular-nums ${winner === 'home' ? 'font-bold text-ink-900' : 'text-ink-500'}`}>
{r.home ?? '—'}
</td>
<td className="py-1.5 text-center text-ink-500">{r.label}</td>
<td className={`py-1.5 text-left tabular-nums ${winner === 'away' ? 'font-bold text-ink-900' : 'text-ink-500'}`}>
{r.away ?? '—'}
</td>
</tr>
)
})}
</tbody>
</table>
</div>
</div>
)
}
/** 近况/H2H 单区块 */
function RecentBlock({ title, rows, side }: { title: string; rows?: TeamRecentMatch[]; side: 'home' | 'away' | 'h2h' }) {
return (
<div>
<h5 className="mb-1.5 text-2xs font-medium text-ink-500">{title}</h5>
{rows && rows.length > 0 ? (
<ul className="space-y-1">
{rows.map((r, i) => {
const date = r.match_date ? new Date(r.match_date).toLocaleDateString('zh-CN', { month: '2-digit', day: '2-digit' }) : '—'
const score = (r.home_goals != null && r.away_goals != null) ? `${r.home_goals}-${r.away_goals}` : 'vs'
const label = side === 'h2h'
? `${r.home_team ?? '?'} ${score} ${r.away_team ?? '?'}`
: `${score}`
return (
<li key={i} className="flex items-center justify-between text-2xs tabular-nums text-ink-600">
<span className="text-ink-400">{date}</span>
<span className="truncate">{label}</span>
</li>
)
})}
</ul>
) : (
<p className="text-2xs text-ink-300"></p>
)}
</div>
)
}
/** 历史预测单行(含专家报告入口) */
function PredictionHistoryRow({ p }: { p: MatchRecentPrediction }) {
const badge = p.status === 'degraded'
? { label: 'degraded', cls: 'text-press' }
: p.settled
? { label: p.correct_1x2 === undefined ? '已结算' : p.correct_1x2 ? '命中' : '未中', cls: p.correct_1x2 ? 'text-ink-900' : 'text-ink-400' }
: { label: p.status === 'success' ? '成功' : p.status, cls: 'text-ink-600' }
const score = (p.pred_home_goals != null && p.pred_away_goals != null)
? `${p.pred_home_goals.toFixed(1)}-${p.pred_away_goals.toFixed(1)}`
: '—'
const alt = (p.alt_pred_home_goals != null && p.alt_pred_away_goals != null)
? `${p.alt_pred_home_goals.toFixed(1)}-${p.alt_pred_away_goals.toFixed(1)}` : null
const hasAgents = p.agent_outputs && p.agent_outputs.length > 0
return (
<div className="border-b border-ink-200 pb-2 last:border-b-0">
<div className="flex items-center justify-between text-xs">
<span className="tabular-nums text-ink-600">
{score} {p.pred_1x2 ? `(${p.pred_1x2})` : ''}
{alt && <span className="ml-1 text-ink-400"> {alt}</span>}
</span>
<span className="flex items-center gap-2">
{p.subjective_confidence != null && (
<span className="text-2xs tabular-nums text-ink-400"> {Math.round(p.subjective_confidence * 100)}%</span>
)}
<span className={`text-2xs ${badge.cls}`}>{badge.label}</span>
</span>
</div>
<div className="mt-0.5 flex items-center justify-between text-2xs text-ink-400">
<span className="truncate">{p.model} · {p.mode} · {p.created_at ? new Date(p.created_at).toLocaleString('zh-CN', { month: '2-digit', day: '2-digit', hour: '2-digit', minute: '2-digit', hour12: false }) : '—'}</span>
{hasAgents && <span className="text-press">{p.agent_outputs!.length} </span>}
</div>
{p.reasoning && (
<p className="mt-1 line-clamp-2 font-serif text-2xs leading-relaxed text-ink-500">{p.reasoning}</p>
)}
</div>
)
}
// 详情懒加载取数在此文件内聚:页面只需切换 expandedId 并缓存结果
export async function loadMatchDetailBundle(
matchId: number,
): Promise<{ detail: MatchDetailOut | null; ctx: MatchContextOut | null }> {
const [d, c] = await Promise.all([
fetchMatchDetail(matchId).catch(() => null),
fetchMatchContext(matchId).catch(() => null),
])
return { detail: d, ctx: c }
}
@@ -0,0 +1,527 @@
/**
* MatchPredictPanel: 预测弹窗全套(过程可视化 / 结果版面 / 专家意见)。
*
* D3: 从 Matches.tsx 拆出,渲染逻辑原样搬迁。对外只导出 PredictModal;
* PredictionPanel 复用 Prediction 的 embedded 模式由弹窗内渲染。
*/
import { useEffect, useState } from 'react'
import TeamSideTag from '../../../components/TeamSideTag'
import type { AgentReport, Match, Prediction } from '../types'
import { AGENT_LABELS, CN_NUM, OUTCOME_LABEL } from '../types'
/** 置信度细线:0~1 数值的低调可视化 */
function Meter({ value }: { value: number }) {
const pct = Math.max(0, Math.min(100, Math.round(value * 100)))
return (
<div className="h-px w-full bg-ink-200" role="presentation">
<div className="h-px bg-press transition-[width] duration-500" style={{ width: `${pct}%` }} />
</div>
)
}
/** home_edge(-1~1,正=利主队)的可视化:以中线为原点的双向细条 */
function EdgeBar({ value }: { value: number }) {
const v = Math.max(-1, Math.min(1, value))
const half = Math.abs(v) * 50
return (
<div className="relative h-px w-full bg-ink-200" role="presentation">
<span className="absolute left-1/2 top-1/2 h-2 w-px -translate-x-1/2 -translate-y-1/2 bg-ink-400" />
<span
className={`absolute top-0 h-px transition-all duration-500 ${v >= 0 ? 'bg-press' : 'bg-ink-600'}`}
style={
v >= 0
? { left: '50%', width: `${half}%` }
: { right: '50%', width: `${half}%` }
}
/>
</div>
)
}
function Spinner({ className = '' }: { className?: string }) {
return (
<svg
viewBox="0 0 20 20"
className={`h-3.5 w-3.5 animate-spin ${className}`}
fill="none"
aria-hidden="true"
>
<circle cx="10" cy="10" r="7.5" stroke="currentColor" strokeWidth="1.5" strokeOpacity="0.25" />
<path d="M17.5 10A7.5 7.5 0 0010 2.5" stroke="currentColor" strokeWidth="1.5" strokeLinecap="round" />
</svg>
)
}
/** 胜平负一行文字:选中的红字加方块标记,未选中的退灰 */
function OutcomeLine({
pick,
confidence,
}: {
pick: string | null
confidence: number | null
}) {
const options = ['1', 'X', '2'] as const
return (
<div>
<div className="flex items-baseline justify-center gap-6 sm:gap-10">
{options.map(o => {
const on = pick === o
return (
<div key={o} className="flex flex-col items-center gap-1">
<span className={`flex items-center gap-1.5 text-sm ${on ? 'font-semibold text-press' : 'text-ink-400'}`}>
{on && <span className="inline-block h-2 w-2 bg-press" aria-hidden="true" />}
{OUTCOME_LABEL[o]}
</span>
{on && confidence !== null && (
<span className="text-2xs tabular-nums text-ink-500">
{Math.round(confidence * 100)}%
</span>
)}
</div>
)
})}
</div>
{pick && confidence !== null && (
<div className="mx-auto mt-3 max-w-xs">
<Meter value={confidence} />
<p className="mt-1 text-center text-2xs text-ink-400">,</p>
</div>
)}
</div>
)
}
/** 预测成本展示:耗时 + token + 限流余量 */
function PredictionCost({ prediction }: { prediction: Prediction }) {
const latency = prediction.latency_ms != null ? `${(prediction.latency_ms / 1000).toFixed(1)}s` : null
const tokens = prediction.prompt_tokens != null || prediction.completion_tokens != null
? `${prediction.prompt_tokens ?? '?'}/${prediction.completion_tokens ?? '?'}`
: null
if (!latency && !tokens && prediction.rate_limit_remaining == null) return null
return (
<div className="border-t border-ink-200 pt-3 text-2xs text-ink-500">
<div className="flex flex-wrap items-center justify-center gap-x-4 gap-y-1">
{latency && (
<span className="inline-flex items-center gap-1">
<span aria-hidden="true" className="opacity-60"></span> {latency}
</span>
)}
{tokens && (
<span className="inline-flex items-center gap-1">
<span aria-hidden="true" className="opacity-60">Tok</span>prompt/completion: {tokens}
</span>
)}
{prediction.rate_limit_remaining != null && prediction.rate_limit_remaining <= 3 && (
<span className="text-press" title="每分钟最多 10 次预测">
: {prediction.rate_limit_remaining}/10()
</span>
)}
</div>
</div>
)
}
/** 预测过程阶段(按时长模拟;结果到达即跳到完成) */
function PredictProgress() {
const [elapsed, setElapsed] = useState(0)
useEffect(() => {
const t = setInterval(() => setElapsed(e => e + 0.5), 500)
return () => clearInterval(t)
}, [])
// 阶段阈值(秒): 切片 → 专家(各路依次点亮) → 终裁
const SLICE_END = 3
const AGENT_START = 4
const AGENT_STEP = 8 // 每路专家约 8s 点亮一路
const AGG_START = AGENT_START + AGENT_STEP * 5
const agents = ['form', 'stats', 'home_away', 'standings', 'h2h']
const phase = elapsed < SLICE_END ? 'slice'
: elapsed < AGG_START ? 'agents' : 'agg'
const pct = Math.min(95, Math.round((elapsed / 70) * 100))
return (
<div className="px-5 py-8 sm:px-8">
{/* 阶段标题 */}
<div className="flex items-center justify-center gap-2">
<Spinner className="text-press" />
<span className="font-serif text-sm font-bold text-ink-900">
{phase === 'slice' && '正在组装比赛数据切片'}
{phase === 'agents' && '五路专家并行分析中'}
{phase === 'agg' && '终裁专家汇总裁定中'}
</span>
<span className="text-2xs tabular-nums text-ink-400">{elapsed.toFixed(0)}s</span>
</div>
{/* 进度条:渐进式,不封顶到 100% */}
<div className="mx-auto mt-5 h-1 w-full max-w-md overflow-hidden bg-ink-100" role="progressbar" aria-valuenow={pct}>
<div
className={`h-full bg-press transition-all duration-500 ${phase === 'agg' ? 'animate-pulse' : ''}`}
style={{ width: `${pct}%` }}
/>
</div>
{/* 专家灯序(多专家模式) */}
<ul className="mx-auto mt-6 max-w-md space-y-1.5">
{agents.map((a, i) => {
const lit = elapsed >= AGENT_START + AGENT_STEP * (i + 1)
const activeNow = !lit && elapsed >= AGENT_START + AGENT_STEP * i
return (
<li
key={a}
className={`flex items-center justify-between border-b border-ink-200 pb-1.5 text-xs transition-colors ${
lit ? 'text-ink-800' : activeNow ? 'text-ink-900' : 'text-ink-300'
}`}
>
<span className="flex items-center gap-2">
<span
aria-hidden="true"
className={`inline-block h-1.5 w-1.5 ${lit ? 'bg-ink-900' : activeNow ? 'bg-press animate-pulse' : 'bg-ink-200'}`}
/>
{AGENT_LABELS[a] ?? a}
</span>
{lit && <span className="text-2xs text-ink-400"> </span>}
{activeNow && <span className="text-2xs text-press"></span>}
</li>
)
})}
</ul>
<p className="mt-6 text-center text-2xs text-ink-400">
, 30-90 ; token,使
</p>
<p className="mt-1 text-center text-2xs text-ink-300">
提示:每分钟限 10 ,
</p>
</div>
)
}
const STATUS_BADGE: Record<string, { label: string; cls: string }> = {
ok: { label: '正常', cls: 'text-ink-500' },
no_data: { label: '无数据', cls: 'text-ink-400' },
error: { label: '调用失败', cls: 'text-press' },
parse_error: { label: '解析失败', cls: 'text-press' },
}
const SUFFICIENCY_LABEL: Record<string, string> = {
high: '充分',
medium: '一般',
low: '偏少',
none: '无',
}
/** 单路专家意见:汉字编号 + 细线行 */
function AgentCard({ report: r, no }: { report: AgentReport; no: string }) {
const badge = STATUS_BADGE[r.status] ?? { label: r.status, cls: 'text-ink-400' }
const inactive = r.status !== 'ok'
return (
<details className="group border-b border-ink-200">
<summary className="flex cursor-pointer list-none items-baseline gap-2.5 px-1 py-3">
<span className="font-serif text-sm text-ink-400">{no}</span>
<span className="text-sm font-medium text-ink-900">{AGENT_LABELS[r.agent] ?? r.agent}</span>
<span className={`text-2xs ${badge.cls}`}>{badge.label}</span>
<span className="ml-auto flex items-baseline gap-3 text-2xs tabular-nums text-ink-500">
{r.status === 'ok' && r.subjective_confidence !== null && (
<span> {Math.round(r.subjective_confidence * 100)}%</span>
)}
{r.status === 'ok' && r.probable_score && (
<span className="font-serif font-bold text-ink-800">{r.probable_score}</span>
)}
<svg viewBox="0 0 20 20" className="h-3 w-3 self-center text-ink-300 transition-transform group-open:rotate-90" fill="currentColor" aria-hidden="true">
<path d="M7.3 5.3a1 1 0 011.4 0l4 4a1 1 0 010 1.4l-4 4a1 1 0 01-1.4-1.4L10.6 10 7.3 6.7a1 1 0 010-1.4z" />
</svg>
</span>
</summary>
<div className="space-y-3 px-1 pb-4 pl-7">
{/* 无数据 / 失败时给出明确说明,避免用户以为是空白 bug */}
{inactive && (
<p className="text-xs leading-relaxed text-ink-500">
{r.status === 'no_data' && '该维度没有可用数据,已跳过 LLM 分析以节省额度(不影响其他专家)。'}
{r.status === 'error' && '该专家调用失败,本次结论未纳入其视角(fail-open 设计,不阻断整体预测)。'}
{r.status === 'parse_error' && '模型输出未通过格式校验,该报告已丢弃。'}
</p>
)}
{!inactive && r.home_edge !== null && (
<div>
<div className="mb-1.5 flex items-baseline justify-between text-2xs">
<span className="text-ink-500"></span>
<span className={`font-semibold tabular-nums ${r.home_edge > 0 ? 'text-press' : r.home_edge < 0 ? 'text-ink-700' : 'text-ink-500'}`}>
{r.home_edge > 0 ? '+' : ''}{r.home_edge.toFixed(2)}
</span>
</div>
<EdgeBar value={r.home_edge} />
<div className="mt-1 flex justify-between text-2xs text-ink-400">
<span></span>
<span></span>
</div>
</div>
)}
{r.analysis && (
<p className="font-serif text-sm leading-loose text-ink-700">{r.analysis}</p>
)}
{r.key_evidence.length > 0 && (
<ul className="space-y-1.5">
{r.key_evidence.map((e, i) => (
<li key={i} className="flex gap-2 text-xs leading-relaxed text-ink-600">
<span className="flex-shrink-0 text-ink-300" aria-hidden="true"></span>
<span>{e}</span>
</li>
))}
</ul>
)}
{r.exp_home_goals !== null && r.exp_away_goals !== null && (
<p className="text-xs text-ink-500">
<span className="font-serif font-bold tabular-nums text-ink-900">{r.exp_home_goals.toFixed(1)} - {r.exp_away_goals.toFixed(1)}</span>
</p>
)}
{!inactive && (
<p className="border-t border-ink-100 pt-2.5 text-2xs text-ink-400">
{SUFFICIENCY_LABEL[r.data_sufficiency] ?? r.data_sufficiency}
<span className="mx-2 text-ink-200">|</span>
<span className="font-mono">{r.model}</span>
{r.latency_ms !== null && <span className="ml-2 tabular-nums">{r.latency_ms}ms</span>}
</p>
)}
</div>
</details>
)
}
function PredictionPanel({
prediction,
match,
embedded = false,
}: {
prediction: Prediction
match: Match
/** 弹窗嵌入模式:弹窗已提供报头,这里省略自带版头 */
embedded?: boolean
}) {
const homeName = match.home_team_zh || match.home_team
const [expertsOpen, setExpertsOpen] = useState(false)
const awayName = match.away_team_zh || match.away_team
const degraded = prediction.status === 'degraded' || prediction.status === 'failed'
const reports = prediction.agent_outputs ?? []
const okReports = reports.filter(r => r.status === 'ok')
return (
<article className={embedded ? 'bg-paper-50' : 'border border-ink-900 bg-paper-50'}>
{!embedded && (
<div className="flex flex-wrap items-baseline justify-between gap-2 border-b border-ink-900 bg-paper-100 px-4 py-2.5 sm:px-5">
<h3 className="flex flex-wrap items-center gap-1.5 font-serif text-sm font-bold text-ink-900">
·
<TeamSideTag side="home" />
{homeName}
<span></span>
<TeamSideTag side="away" />
{awayName}
</h3>
<span className="text-2xs tabular-nums text-ink-500">
{prediction.provider} / {prediction.model}
{prediction.latency_ms !== null && ` · ${(prediction.latency_ms / 1000).toFixed(1)}s`}
</span>
</div>
)}
<div className="space-y-7 px-4 py-6 sm:px-5">
{/* ── degraded / failed 态:醒目警示 + 原因,不展示虚假比分 ── */}
{degraded && (
<div className="border-l-2 border-press bg-press-wash/40 px-4 py-3">
<p className="font-serif text-sm font-bold text-press-dark">
{prediction.status === 'failed' ? '预测失败' : '预测降级(degraded)'}
</p>
<p className="mt-1.5 whitespace-pre-wrap text-xs leading-relaxed text-ink-600">
{prediction.reasoning || '所有专家均无有效数据或调用失败,无法生成可靠比分。'}
</p>
</div>
)}
{/* ── 主结论(仅 success 展示) ── */}
{!degraded && (
<>
<div className="text-center">
<p className="font-serif text-5xl font-bold tabular-nums leading-none text-ink-900 sm:text-6xl">
{prediction.pred_home_goals ?? '-'}
<span className="mx-3 font-normal text-ink-300">:</span>
{prediction.pred_away_goals ?? '-'}
</p>
<p className="mt-3 text-2xs tracking-[0.5em] text-ink-400"></p>
{prediction.alt_pred_home_goals != null && prediction.alt_pred_away_goals != null && (
<p className="mt-2 text-2xs tabular-nums text-ink-400">
{' '}
<span className="font-serif text-sm font-bold tabular-nums text-ink-600">
{prediction.alt_pred_home_goals}<span className="mx-0.5 font-normal text-ink-300">:</span>{prediction.alt_pred_away_goals}
</span>
</p>
)}
</div>
<div className="border-y border-ink-200 py-4">
<OutcomeLine pick={prediction.pred_1x2} confidence={prediction.subjective_confidence} />
</div>
</>
)}
{/* ── 成本信息(耗时 + token + 限流余量) ── */}
{!degraded && (
<PredictionCost prediction={prediction} />
)}
{/* ── 元信息 ── */}
<p className="text-center text-2xs text-ink-500">
`多专家模式 · ${okReports.length}/${reports.length} 路有效`
{prediction.prompt_version && ` · prompt ${prediction.prompt_version}`}
</p>
{/* ── 终裁/降级说明意见 ── */}
{prediction.reasoning && degraded && (
<section>
<h4 className="section-head mb-2"></h4>
<blockquote className="border-l-2 border-press pl-4">
<p className="whitespace-pre-wrap font-serif text-sm leading-loose text-ink-700">{prediction.reasoning}</p>
</blockquote>
</section>
)}
{/* ── 专家意见(多模式):可折叠 + 状态摘要 + 权重条形图 ── */}
{reports.length > 0 && (
<section>
<button
onClick={() => setExpertsOpen(o => !o)}
className="flex w-full items-center justify-between border-b border-ink-200 pb-2 text-left"
>
<span className="section-head mb-0">({okReports.length}/{reports.length} )</span>
<span className="text-2xs text-ink-400">{expertsOpen ? '收起' : '展开'}</span>
</button>
{/* 权重条形图(仅 success 且有权重时显示) */}
{!degraded && prediction.agent_weights && Object.keys(prediction.agent_weights).length > 0 && (
<div className="mt-3 space-y-1.5">
<span className="text-2xs text-ink-500"></span>
{Object.entries(prediction.agent_weights)
.sort((a, b) => b[1] - a[1])
.map(([k, v]) => (
<div key={k} className="grid grid-cols-[96px_minmax(0,1fr)_40px] items-center gap-2">
<span className="truncate text-2xs text-ink-500">{AGENT_LABELS[k] ?? k}</span>
<div className="h-1.5 bg-paper-100">
<div className="h-full bg-press" style={{ width: `${Math.round(v * 100)}%` }} />
</div>
<span className="text-right text-2xs tabular-nums text-ink-500">{Math.round(v * 100)}%</span>
</div>
))}
</div>
)}
{expertsOpen && (
<div className="mt-2">
{reports.map((r, i) => (
<AgentCard key={r.agent} report={r} no={CN_NUM[i] ?? String(i + 1)} />
))}
</div>
)}
</section>
)}
{/* ── 终裁意见(success) ── */}
{prediction.reasoning && !degraded && (
<section>
<h4 className="section-head mb-3"></h4>
<blockquote className="border-l-2 border-press pl-4">
<p className="whitespace-pre-wrap font-serif text-sm leading-loose text-ink-700">{prediction.reasoning}</p>
</blockquote>
</section>
)}
</div>
</article>
)
}
/** 预测弹窗:进行中显示过程可视化,完成后显示预测版,失败显示原因 */
export function PredictModal({
match,
predicting,
prediction,
error,
onClose,
}: {
match: Match
predicting: boolean
prediction: Prediction | null
error: string | null
onClose: () => void
}) {
const homeName = match.home_team_zh || match.home_team
const awayName = match.away_team_zh || match.away_team
useEffect(() => {
const h = (e: KeyboardEvent) => {
if (e.key === 'Escape') onClose()
}
document.addEventListener('keydown', h)
return () => document.removeEventListener('keydown', h)
}, [onClose])
return (
<div
className="fixed inset-0 z-50 flex items-start justify-center overflow-y-auto bg-ink-900/50 p-4 sm:items-center"
role="dialog"
aria-modal="true"
aria-label={`预测 ${homeName}${awayName}`}
onClick={e => {
if (e.target === e.currentTarget) onClose()
}}
>
<div className="relative flex max-h-[92vh] w-full max-w-2xl flex-col overflow-hidden bg-paper-50 shadow-2xl">
{/* 弹窗报头 */}
<div className="flex flex-shrink-0 items-center justify-between border-b border-ink-900 bg-paper-100 px-4 py-2.5 sm:px-5">
<h3 className="flex flex-wrap items-center gap-1.5 font-serif text-sm font-bold text-ink-900">
·
<TeamSideTag side="home" />
{homeName}
<span></span>
<TeamSideTag side="away" />
{awayName}
</h3>
<button
onClick={onClose}
className="flex h-11 w-11 flex-shrink-0 items-center justify-center text-ink-400 transition-colors hover:text-ink-900"
aria-label="关闭"
>
<svg className="h-4 w-4" viewBox="0 0 20 20" fill="currentColor" aria-hidden="true">
<path d="M6.3 5.3a1 1 0 011.4 0L10 7.6l2.3-2.3a1 1 0 111.4 1.4L11.4 9l2.3 2.3a1 1 0 01-1.4 1.4L10 10.4l-2.3 2.3a1 1 0 01-1.4-1.4L8.6 9 6.3 6.7a1 1 0 010-1.4z" />
</svg>
</button>
</div>
{/* 弹窗体(小屏可滚动) */}
<div className="flex-1 overflow-y-auto">
{predicting ? (
<PredictProgress />
) : error ? (
<div className="px-5 py-10 text-center sm:px-8">
<p className="font-serif text-sm font-bold text-press"></p>
<p className="mx-auto mt-3 max-w-md whitespace-pre-wrap text-left text-xs leading-relaxed text-ink-600">
{error}
</p>
<button onClick={onClose} className="btn btn-sm mt-6"></button>
</div>
) : prediction ? (
<PredictionPanel prediction={prediction} match={match} embedded />
) : null}
</div>
</div>
</div>
)
}
@@ -0,0 +1,93 @@
/**
* useMatchPredict: 预测流程状态机(发起/进行中/结果/失败/关闭中止)。
*
* D3: 从 Matches.tsx 拆出。语义不变:
* - 连点防护:同一场比赛预测中再次点击直接忽略
* - 竞态防护:递增序号,过期响应丢弃
* - 关闭弹窗 = 中止在途请求 + 序号失效(catch/then 不再写入)
* - 5 分钟超时,与 nginx 代理 300s 对齐
*/
import { useRef, useState } from 'react'
import { http } from '../../../lib/http'
import type { Match, Prediction } from '../types'
/** 把后端/网络错误翻译成用户可读文案 */
function readablePredictError(e: unknown): string {
if (e instanceof Error) {
const m = e.message
if (/429/.test(m)) {
// 429 来自后端限流(每分钟 10 次),非上游 LLM
return '操作过于频繁:每分钟最多 10 次预测。为保护 LLM 额度,请稍后再试。'
}
if (/502/.test(m)) return 'LLM 服务暂时不可用(502),请稍后重试'
if (/402|Payment Required|额度|余额/.test(m)) return 'LLM 额度不足(402),请检查 API Key 余额'
if (/400|已完赛/.test(m)) return '该比赛已完赛,不再支持预测'
if (/409|已结算/.test(m)) return '该预测已结算,不能重新预测'
if (/timeout|超时|timed out/i.test(m)) return '请求超时,请稍后重试'
return m
}
return String(e)
}
interface UseMatchPredictOptions {
/** 共享 error state(拆分前列表与预测共用同一个 error,行为保持一致) */
onError: (msg: string | null) => void
}
export function useMatchPredict({ onError }: UseMatchPredictOptions) {
const [predictingId, setPredictingId] = useState<number | null>(null)
const [prediction, setPrediction] = useState<Prediction | null>(null)
const [predictionFor, setPredictionFor] = useState<Match | null>(null)
const predictSeq = useRef(0)
// 预测请求控制器:关闭弹窗时中止
const predictAbort = useRef<AbortController | null>(null)
function closePredict() {
predictAbort.current?.abort()
predictSeq.current++ // 令中止请求的 catch/then 全部失效,不再写入错误
setPredictingId(null)
setPrediction(null)
setPredictionFor(null)
onError(null)
}
const predict = async (m: Match) => {
// 防连点:若该场比赛已在预测中,直接忽略
if (predictingId === m.id) return
const seq = ++predictSeq.current
setPredictingId(m.id)
onError(null)
setPrediction(null)
setPredictionFor(m)
// LLM 多专家预测耗时可达数分钟,给足超时(与 nginx 代理 300s 对齐)
const controller = new AbortController()
predictAbort.current = controller
const timer = setTimeout(() => controller.abort(), 300_000)
try {
const data = await http.post<Prediction>('/predict', { match_id: m.id, mode: 'multi' }, {
timeoutMs: 300_000,
signal: controller.signal,
})
if (seq !== predictSeq.current) return
setPrediction(data)
} catch (e) {
if (seq !== predictSeq.current) return
onError(
e instanceof DOMException && e.name === 'AbortError'
? '预测超时(5 分钟),请稍后重试'
: readablePredictError(e),
)
} finally {
clearTimeout(timer)
if (seq === predictSeq.current) setPredictingId(null)
}
}
return {
predictingId,
prediction,
predictionFor,
predict,
closePredict,
}
}
@@ -0,0 +1,94 @@
/**
* useMatchesList: 赛程列表数据获取(筛选/游标分页/进行中比赛)。
*
* D3: 从 Matches.tsx 拆出。竞态防护语义不变 —— 递增序号只认最后一次请求;
* loadMore 不自增序号(切换筛选才自增,翻页跟随当前序列)。
*/
import { useCallback, useEffect, useRef, useState } from 'react'
import { http } from '../../../lib/http'
import type { Match } from '../types'
interface UseMatchesListOptions {
/** 共享 error state(拆分前列表与预测共用同一个 error,行为保持一致) */
onError: (msg: string | null) => void
}
export function useMatchesList({ onError }: UseMatchesListOptions) {
const [league, setLeague] = useState('E0')
const [status, setStatus] = useState('scheduled')
const [matches, setMatches] = useState<Match[]>([])
const [nextCursor, setNextCursor] = useState<string | null>(null)
const [loadingMore, setLoadingMore] = useState(false)
const [loading, setLoading] = useState(false)
const [showAllUpcoming, setShowAllUpcoming] = useState(false) // 默认仅展示未来 3 天;true 展开全部
const [liveMatches, setLiveMatches] = useState<Match[]>([]) // 进行中比赛(顶部独立区块)
// 请求竞态防护:切换联赛/状态很快时,先发的慢请求可能后返回,
// 把旧结果覆盖到新筛选上。用递增序号只认最后一次请求的响应。
const loadSeq = useRef(0)
const load = useCallback(async () => {
const seq = ++loadSeq.current
setLoading(true)
setLoadingMore(false)
setShowAllUpcoming(false) // 切换筛选重置为「未来 3 天」视图
onError(null)
try {
const params = new URLSearchParams({ league, status, limit: '50' })
const data = await http.get<{ items: Match[]; next_cursor: string | null }>(`/matches?${params}`)
if (seq !== loadSeq.current) return // 已有更新的请求,丢弃本次结果
setMatches(data.items)
setNextCursor(data.next_cursor ?? null)
} catch (e) {
if (seq !== loadSeq.current) return
onError(e instanceof Error ? e.message : String(e))
} finally {
if (seq === loadSeq.current) setLoading(false)
}
}, [league, status, onError])
// 加载下一页(游标分页)
const loadMore = async () => {
if (!nextCursor || loadingMore) return
const seq = loadSeq.current // 不做自增:切换筛选会自增,这里只跟随当前序列
setLoadingMore(true)
try {
const params = new URLSearchParams({ league, status, limit: '50', cursor: nextCursor })
const data = await http.get<{ items: Match[]; next_cursor: string | null }>(`/matches?${params}`)
if (seq !== loadSeq.current) return
setMatches(prev => [...prev, ...data.items])
setNextCursor(data.next_cursor ?? null)
} catch (e) {
if (seq !== loadSeq.current) return
onError(e instanceof Error ? e.message : String(e))
} finally {
if (seq === loadSeq.current) setLoadingMore(false)
}
}
// 加载进行中比赛(顶部独立区块)
const loadLive = useCallback(async () => {
try {
const params = new URLSearchParams({ league, status: 'in_play', limit: '20' })
const data = await http.get<{ items: Match[] }>(`/matches?${params}`)
setLiveMatches(data.items ?? [])
} catch {
/* ignore:进行中非核心功能 */
}
}, [league])
useEffect(() => { load(); loadLive() }, [load, loadLive])
return {
league, setLeague,
status, setStatus,
matches,
nextCursor,
loading,
loadingMore,
showAllUpcoming, setShowAllUpcoming,
liveMatches,
load,
loadMore,
}
}
+98
View File
@@ -0,0 +1,98 @@
/**
* Matches 页面族共享类型与常量。
*
* D3(工程债): Matches.tsx 原本 ~1400 行,类型/常量/预测面板/详情面板/
* 列表逻辑全部内联。本文件是拆分后的共享层 —— 只放数据契约与纯常量,
* 不含 React 组件。
*/
/** /matches 列表项(公开接口) */
export interface Match {
id: number
league_code: string | null
season: string | null
home_team: string
away_team: string
home_team_zh: string | null
away_team_zh: string | null
match_date: string
match_status: string
home_goals: number | null
away_goals: number | null
match_stage: string | null
home_xg: number | null
away_xg: number | null
}
/** POST /predict 响应 */
export interface Prediction {
prediction_id: number
provider: string
model: string
prompt_version: string | null
mode: string
pred_home_goals: number | null
pred_away_goals: number | null
alt_pred_home_goals: number | null
alt_pred_away_goals: number | null
pred_1x2: string | null
subjective_confidence: number | null
reasoning: string | null
status: string
agent_outputs: AgentReport[] | null
agent_weights: Record<string, number> | null
context: string
latency_ms: number | null
prompt_tokens: number | null
completion_tokens: number | null
rate_limit_remaining: number | null
}
/** 多专家单路报告 */
export interface AgentReport {
agent: string
status: string
data_sufficiency: string
analysis: string
home_edge: number | null
subjective_confidence: number | null
key_evidence: string[]
exp_home_goals: number | null
exp_away_goals: number | null
probable_score: string | null
model: string
latency_ms: number | null
}
export const AGENT_LABELS: Record<string, string> = {
h2h: '历史交锋分析专家',
form: '近期状态分析专家',
stats: '攻防数据分析专家',
home_away: '主客因素分析专家',
standings: '联赛排名分析专家',
}
export const LEAGUES = [
{ code: 'E0', name: '英超' },
{ code: 'SP1', name: '西甲' },
{ code: 'D1', name: '德甲' },
{ code: 'I1', name: '意甲' },
{ code: 'F1', name: '法甲' },
]
/** 汉字编号,给专家意见排版用 */
export const CN_NUM = ['一', '二', '三', '四', '五', '六', '七', '八']
export const STATUS_META: Record<string, { label: string; cls: string }> = {
finished: { label: '已完赛', cls: 'text-ink-400' },
scheduled: { label: '未开赛', cls: 'text-ink-600' },
// 键与 normalize.py 的 VALID_STATUS 对齐: 库里存的是 in_play(上游 live 被归一化),不存在 'live' 状态
in_play: { label: '进行中', cls: 'text-press font-medium' },
paused: { label: '暂停', cls: 'text-press font-medium' },
postponed: { label: '延期', cls: 'text-ink-400' },
cancelled: { label: '取消', cls: 'text-ink-400' },
suspended: { label: '中止', cls: 'text-ink-400' },
}
/** 1x2 → 中文标签 */
export const OUTCOME_LABEL: Record<string, string> = { '1': '主胜', X: '平局', '2': '客胜' }
+114
View File
@@ -0,0 +1,114 @@
/**
* Matches 页面族共享的原子 UI 小件(无业务状态)。
*
* D3: 从 Matches.tsx 内联定义上移到模块级 —— Switch 原先定义在组件函数
* 体内(每次渲染重建组件对象),它没有内部 state,提升后渲染结果一致。
*/
import type { Match } from './types'
export function Spinner({ className = '' }: { className?: string }) {
return (
<svg
viewBox="0 0 20 20"
className={`h-3.5 w-3.5 animate-spin ${className}`}
fill="none"
aria-hidden="true"
>
<circle cx="10" cy="10" r="7.5" stroke="currentColor" strokeWidth="1.5" strokeOpacity="0.25" />
<path d="M17.5 10A7.5 7.5 0 0010 2.5" stroke="currentColor" strokeWidth="1.5" strokeLinecap="round" />
</svg>
)
}
/** 骨架占位行:低调脉动灰块 */
export function SkeletonRows({ n = 4 }: { n?: number }) {
return (
<>
{Array.from({ length: n }).map((_, i) => (
<div key={i} className="flex items-center gap-4 border-b border-ink-200 px-1 py-3.5">
<div className="skeleton h-3 w-16" />
<div className="skeleton h-3 flex-1" />
<div className="skeleton h-3 w-10" />
<div className="skeleton h-3 flex-1" />
<div className="skeleton h-3 w-16" />
</div>
))}
</>
)
}
/** 状态/模式一组的文字切换 */
export function Switch({ value, onChange, items }: {
value: string
onChange: (v: string) => void
items: { v: string; label: string; title?: string }[]
}) {
return (
<span className="inline-flex items-center gap-2.5">
{items.map((it, i) => (
<span key={it.v} className="inline-flex items-center gap-2.5">
{i > 0 && <span className="text-ink-300" aria-hidden="true">/</span>}
<button
onClick={() => onChange(it.v)}
title={it.title}
className={`relative tab ${value === it.v ? 'tab-on' : ''} text-xs`}
>
{it.label}
</button>
</span>
))}
</span>
)
}
/** 日期分组头显示:今日/明天/周几 · 年月日 */
export function formatDateHeader(dateKey: string): string {
if (!dateKey) return '未开赛'
const d = new Date(dateKey + 'T00:00:00')
if (isNaN(d.getTime())) return dateKey
const today = new Date()
const todayKey = `${today.getFullYear()}-${String(today.getMonth() + 1).padStart(2, '0')}-${String(today.getDate()).padStart(2, '0')}`
const tmr = new Date(today)
tmr.setDate(tmr.getDate() + 1)
const tmrKey = `${tmr.getFullYear()}-${String(tmr.getMonth() + 1).padStart(2, '0')}-${String(tmr.getDate()).padStart(2, '0')}`
const weekday = ['周日', '周一', '周二', '周三', '周四', '周五', '周六'][d.getDay()]
if (dateKey === todayKey) return `今日 ${weekday}`
if (dateKey === tmrKey) return `明日 ${weekday}`
return `${d.getMonth() + 1}${d.getDate()}${weekday}`
}
/** UTC ISO → 本地日期 YYYY-MM-DD(用于分组) */
export function toLocalDateKey(iso: string): string {
const d = new Date(iso)
return `${d.getFullYear()}-${String(d.getMonth() + 1).padStart(2, '0')}-${String(d.getDate()).padStart(2, '0')}`
}
/** 按本地日期分组(非 UTC),保持时间序 */
export function groupByDate(list: Match[]): Array<[string, Match[]]> {
const map = new Map<string, Match[]>()
for (const m of list) {
const key = toLocalDateKey(m.match_date)
const arr = map.get(key)
if (arr) arr.push(m)
else map.set(key, [m])
}
return [...map.entries()]
}
/** 日期 key 辅助:YYYY-MM-DD(本地时区) */
function dateKey(d: Date): string {
return `${d.getFullYear()}-${String(d.getMonth() + 1).padStart(2, '0')}-${String(d.getDate()).padStart(2, '0')}`
}
/** 未来 3 天窗口:今天 00:00 → 第 3 天 00:00(即今天/明天/后天) */
function addDays(d: Date, n: number): string {
const x = new Date(d)
x.setFullYear(x.getFullYear(), x.getMonth(), x.getDate() + n)
return dateKey(x)
}
/** 比赛是否在未来 3 天内(用于默认视图过滤) */
export function withinNext3Days(matchDate: string): boolean {
const key = toLocalDateKey(matchDate)
return key >= dateKey(new Date()) && key < addDays(new Date(), 3)
}
+10
View File
@@ -30,6 +30,16 @@ async def lifespan(app: FastAPI) -> AsyncIterator[None]:
await ensure_admin_password_hashed() # .env 明文密码 → scrypt 哈希(幂等) await ensure_admin_password_hashed() # .env 明文密码 → scrypt 哈希(幂等)
await assert_security_on_startup() # 启动安全校验(生产拒绝/开发警告) await assert_security_on_startup() # 启动安全校验(生产拒绝/开发警告)
# D7(工程债): 进程内限流(_RateLimiter)与 KeyRing 均为单进程状态;
# 多 worker 部署时各进程独立计数,限流阈值会按 worker 数放大、KeyRing 不共享。
# 生产环境应将限流前置到 Nginx/网关,或以单 worker 运行(见 src/api/deps.py 注释)。
# 此处仅提醒一次,不阻断启动,也不引入 Redis 等外部依赖。
if settings.APP_ENV == "production":
logger.warning(
"APP_ENV=production: 进程内 rate-limit 与 KeyRing 仅单进程有效;"
"多 worker 部署请将限流前置到 Nginx/网关,或以单 worker 运行"
)
# 注册默认定时任务(如果数据库中没有) # 注册默认定时任务(如果数据库中没有)
from src.db.base import AsyncSessionLocal from src.db.base import AsyncSessionLocal
from sqlalchemy import select from sqlalchemy import select
+31 -34
View File
@@ -63,7 +63,8 @@ async def predict(req: PredictRequest, request: Request):
logger.exception("predict unexpected error") logger.exception("predict unexpected error")
raise HTTPException(500, "预测失败,请查看服务器日志") raise HTTPException(500, "预测失败,请查看服务器日志")
# baseline 模式:结果已是 dict,需独立落库(prediction_id) # D2: 三种模式统一返回 PredictResult —— 字段映射单一化,无 dict 分支。
# 仅 baseline 的 prediction_id 需要在此落库补齐(服务层不落库)。
if req.mode == "baseline": if req.mode == "baseline":
prediction_id = await _persist_baseline(req.match_id, result) prediction_id = await _persist_baseline(req.match_id, result)
else: else:
@@ -73,38 +74,34 @@ async def predict(req: PredictRequest, request: Request):
logger.info( logger.info(
"预测完成 match=%s mode=%s pred=%s:%s (%s)", "预测完成 match=%s mode=%s pred=%s:%s (%s)",
req.match_id, req.mode, req.match_id, req.mode,
result.get("pred_home_goals") if isinstance(result, dict) else result.pred_home_goals, result.pred_home_goals, result.pred_away_goals, result.pred_1x2,
result.get("pred_away_goals") if isinstance(result, dict) else result.pred_away_goals,
result.get("pred_1x2") if isinstance(result, dict) else result.pred_1x2,
) )
result_dict = result if isinstance(result, dict) else None
return PredictOut( return PredictOut(
prediction_id=prediction_id, prediction_id=prediction_id,
provider=result.get("provider") if result_dict else result.provider, provider=result.provider,
model=result.get("model") if result_dict else result.model, model=result.model,
prompt_version=result.get("prompt_version") if result_dict else getattr(result, "prompt_version", None), prompt_version=result.prompt_version,
mode=req.mode, mode=req.mode,
pred_home_goals=result.get("pred_home_goals") if result_dict else result.pred_home_goals, pred_home_goals=result.pred_home_goals,
pred_away_goals=result.get("pred_away_goals") if result_dict else result.pred_away_goals, pred_away_goals=result.pred_away_goals,
alt_pred_home_goals=result.get("alt_pred_home_goals") if result_dict else result.alt_pred_home_goals, alt_pred_home_goals=result.alt_pred_home_goals,
alt_pred_away_goals=result.get("alt_pred_away_goals") if result_dict else result.alt_pred_away_goals, alt_pred_away_goals=result.alt_pred_away_goals,
pred_1x2=result.get("pred_1x2") if result_dict else result.pred_1x2, pred_1x2=result.pred_1x2,
subjective_confidence=result.get("subjective_confidence") if result_dict else result.subjective_confidence, subjective_confidence=result.subjective_confidence,
reasoning=result.get("reasoning") if result_dict else result.reasoning, reasoning=result.reasoning,
status=result.get("status", "success") if result_dict else getattr(result, "status", "success"), status=result.status,
agent_outputs=result.get("agent_outputs") if result_dict else getattr(result, "agent_outputs", None), agent_outputs=result.agent_outputs,
agent_weights=result.get("agent_weights") if result_dict else getattr(result, "agent_weights", None), agent_weights=result.agent_weights,
context=result.get("context", "") if result_dict else result.context, context=result.context,
latency_ms=result.get("latency_ms", 0) if result_dict else result.latency_ms, latency_ms=result.latency_ms,
prompt_tokens=result.get("prompt_tokens") if result_dict else getattr(result, "prompt_tokens", None), prompt_tokens=result.prompt_tokens,
completion_tokens=result.get("completion_tokens") if result_dict else getattr(result, "completion_tokens", None), completion_tokens=result.completion_tokens,
rate_limit_remaining=get_predict_rate_limit_remaining(request), rate_limit_remaining=get_predict_rate_limit_remaining(request),
) )
async def _persist_baseline(match_id: int, baseline: dict) -> int: async def _persist_baseline(match_id: int, baseline: PredictResult) -> int:
"""将基线预测结果写入 prediction 表,复用 upsert 语义。""" """将基线预测结果写入 prediction 表,复用 upsert 语义。"""
from src.db.unit_of_work import get_uow from src.db.unit_of_work import get_uow
from src.llm.predict import _upsert_prediction from src.llm.predict import _upsert_prediction
@@ -118,16 +115,16 @@ async def _persist_baseline(match_id: int, baseline: dict) -> int:
mode="baseline", mode="baseline",
run_type="live", # baseline 是 live 预测的变体,符合 ck_run_type_enum run_type="live", # baseline 是 live 预测的变体,符合 ck_run_type_enum
values={ values={
"prompt_version": "baseline_v1", "prompt_version": baseline.prompt_version,
"prompt_tokens": 0, "prompt_tokens": baseline.prompt_tokens or 0,
"completion_tokens": 0, "completion_tokens": baseline.completion_tokens or 0,
"latency_ms": 0, "latency_ms": baseline.latency_ms or 0,
"pred_home_goals": baseline["pred_home_goals"], "pred_home_goals": baseline.pred_home_goals,
"pred_away_goals": baseline["pred_away_goals"], "pred_away_goals": baseline.pred_away_goals,
"pred_1x2": baseline["pred_1x2"], "pred_1x2": baseline.pred_1x2,
"subjective_confidence": baseline["subjective_confidence"], "subjective_confidence": baseline.subjective_confidence,
"reasoning": baseline["reasoning"], "reasoning": baseline.reasoning,
"raw_response": baseline.get("raw", baseline), "raw_response": baseline.raw,
"status": "success", "status": "success",
}, },
) )
+157 -51
View File
@@ -5,7 +5,9 @@
2. standings— 联赛积分榜快照(/leagues/{id}/standings/) 2. standings— 联赛积分榜快照(/leagues/{id}/standings/)
3. stats — 已完赛比赛详细统计回填(/events/{id}/stats/) 3. stats — 已完赛比赛详细统计回填(/events/{id}/stats/)
使用 Repository 模式进行数据访问,不直接控制事务(由调用方 UnitOfWork 控制)。 D4(工程债): Team/League/Match 的查找/创建经 Repository 层(src/db/repositories.py),
本模块不直接控制事务(commit/rollback 由调用方 UnitOfWork 控制,这里只 flush)。
Standing/RawEvent/Lineage 等管线内私有读写仍在本模块内实现,不强行 Repository 化。
""" """
from __future__ import annotations from __future__ import annotations
@@ -16,7 +18,6 @@ from collections.abc import Iterable
from datetime import datetime, timedelta, timezone from datetime import datetime, timedelta, timezone
from sqlalchemy import select from sqlalchemy import select
from sqlalchemy.orm import selectinload
import httpx import httpx
@@ -27,7 +28,8 @@ from src.data.key_ring import _mask, get_key_ring
from src.data.normalize import normalize_bzzoiro from src.data.normalize import normalize_bzzoiro
from src.data.team_names_zh import zh_name from src.data.team_names_zh import zh_name
from src.data.sources import register from src.data.sources import register
from src.db.models import League, Match, MatchStats, Standing, Team, RawEvent, IngestFailure, DataLineage from src.db.models import Match, MatchStats, Standing, Team, RawEvent, IngestFailure, DataLineage
from src.db.repositories import LeagueRepository, MatchRepository, TeamRepository
logger = logging.getLogger(__name__) logger = logging.getLogger(__name__)
@@ -200,16 +202,22 @@ class BzzoiroSource:
# 单联赛抓取失败隔离:记录错误后继续其余联赛,不拖垮整批 # 单联赛抓取失败隔离:记录错误后继续其余联赛,不拖垮整批
logger.exception("bzzoiro fetch failed for %s", code) logger.exception("bzzoiro fetch failed for %s", code)
league_r["errors"].append(f"fetch failed: {e}") league_r["errors"].append(f"fetch failed: {e}")
await _safe_write_ingest_failure(
db,
entity_type="events",
source_record_id=None,
error=e,
raw_payload={"league": code, "status": status, "date_from": date_from, "date_to": date_to},
)
result["leagues"][code] = league_r result["leagues"][code] = league_r
continue continue
# 获取或创建联赛 # D4: 联赛查找/创建经 LeagueRepository(事务仍由调用方 UoW 提交)
stmt = select(League).where(League.code == code) league = await LeagueRepository(db).get_or_create(
league = (await db.execute(stmt)).scalar_one_or_none() code, LEAGUE_NAMES.get(code, code), LEAGUE_COUNTRIES.get(code)
if league is None: )
league = League(code=code, name=LEAGUE_NAMES.get(code, code), country=LEAGUE_COUNTRIES.get(code)) team_r = TeamRepository(db)
db.add(league) match_r = MatchRepository(db)
await db.flush()
# === 批量优化: 预加载球队和已有比赛到内存 === # === 批量优化: 预加载球队和已有比赛到内存 ===
team_name_to_id: dict[str, int] = {} team_name_to_id: dict[str, int] = {}
@@ -235,9 +243,10 @@ class BzzoiroSource:
all_team_names.add(nm.away_team) all_team_names.add(nm.away_team)
if all_team_names: if all_team_names:
stmt = select(Team).where(Team.name.in_(all_team_names)) team_name_to_id = {
teams = (await db.execute(stmt)).scalars().all() name: t.id
team_name_to_id = {t.name: t.id for t in teams} for name, t in (await team_r.get_all_by_names(list(all_team_names))).items()
}
# P1-2: 按需加载,只加载 raw_events 涉及日期范围的比赛(加 30 天缓冲) # P1-2: 按需加载,只加载 raw_events 涉及日期范围的比赛(加 30 天缓冲)
# 避免加载联赛全部历史比赛到内存(多赛季采集时内存溢出) # 避免加载联赛全部历史比赛到内存(多赛季采集时内存溢出)
@@ -247,33 +256,34 @@ class BzzoiroSource:
if dates: if dates:
min_dt = min(dates) - timedelta(days=30) min_dt = min(dates) - timedelta(days=30)
max_dt = max(dates) + timedelta(days=30) max_dt = max(dates) + timedelta(days=30)
stmt = ( matches_in_range = await match_r.find_by_league_and_date_range(
select(Match) league.id, min_dt, max_dt
.where(Match.league_id == league.id)
.where(Match.match_date >= min_dt)
.where(Match.match_date <= max_dt)
) )
existing_matches = { existing_matches = {
_match_key(m.home_team_id, m.away_team_id, m.match_date_date): m _match_key(m.home_team_id, m.away_team_id, m.match_date_date): m
for m in (await db.execute(stmt)).scalars() for m in matches_in_range
} }
# else: existing_matches 保持空 dict(全量新比赛) # else: existing_matches 保持空 dict(全量新比赛)
# D1: Bronze 层批次信息(每联赛每批次一个 batch_id;seen 防同批重复写入)
now = datetime.now(timezone.utc)
bronze_batch_id = f"bzzoiro-events-{code}-{now:%Y%m%d%H%M%S}"
bronze_written: set[str] = set()
for nm, raw in normalized_matches: for nm, raw in normalized_matches:
# 球队: 内存查找 + 按需创建 # D1: RawEvent 幂等键(上游 id 或合成键),插入/变更更新共用
record_id = _events_record_id(code, nm, raw)
# 球队: 内存查找 + 按需创建(D4: 经 TeamRepository)
home_team_id = team_name_to_id.get(nm.home_team) home_team_id = team_name_to_id.get(nm.home_team)
if home_team_id is None: if home_team_id is None:
home = Team(name=nm.home_team, name_zh=zh_name(nm.home_team)) home = await team_r.get_or_create(nm.home_team, name_zh=zh_name(nm.home_team))
db.add(home)
await db.flush()
home_team_id = home.id home_team_id = home.id
team_name_to_id[nm.home_team] = home_team_id team_name_to_id[nm.home_team] = home_team_id
away_team_id = team_name_to_id.get(nm.away_team) away_team_id = team_name_to_id.get(nm.away_team)
if away_team_id is None: if away_team_id is None:
away = Team(name=nm.away_team, name_zh=zh_name(nm.away_team)) away = await team_r.get_or_create(nm.away_team, name_zh=zh_name(nm.away_team))
db.add(away)
await db.flush()
away_team_id = away.id away_team_id = away.id
team_name_to_id[nm.away_team] = away_team_id team_name_to_id[nm.away_team] = away_team_id
@@ -303,6 +313,18 @@ class BzzoiroSource:
# 统计字段不在 /events/ 载荷中(单独由 stats 管线回填), # 统计字段不在 /events/ 载荷中(单独由 stats 管线回填),
# 此处不再创建 MatchStats。 # 此处不再创建 MatchStats。
league_r["inserted"] += 1 league_r["inserted"] += 1
# D1: 成功插入 → 补写 Bronze 层(原始载荷 + 血缘)
if record_id not in bronze_written:
bronze_written.add(record_id)
await _write_events_bronze(
db,
source_record_id=record_id,
raw_payload=raw,
target_match_id=m.id,
league_code=code,
match_status=nm.match_status,
batch_id=bronze_batch_id,
)
else: else:
# 已有比赛: 直接从内存获取对象更新(无需再查询) # 已有比赛: 直接从内存获取对象更新(无需再查询)
changed = False changed = False
@@ -325,6 +347,18 @@ class BzzoiroSource:
changed = True changed = True
if changed: if changed:
league_r["updated"] += 1 league_r["updated"] += 1
# D1: 变更更新 → 补写血缘(RawEvent 幂等键不变,重复采集自动跳过)
if record_id not in bronze_written:
bronze_written.add(record_id)
await _write_events_bronze(
db,
source_record_id=record_id,
raw_payload=raw,
target_match_id=existing_match.id,
league_code=code,
match_status=nm.match_status,
batch_id=bronze_batch_id,
)
# 注意: 不在此处 commit,由调用方 UnitOfWork 控制事务 # 注意: 不在此处 commit,由调用方 UnitOfWork 控制事务
result["leagues"][code] = league_r result["leagues"][code] = league_r
@@ -367,6 +401,31 @@ async def _write_ingest_failure(db, source_system: str, entity_type: str, source
)) ))
async def _safe_write_ingest_failure(
db,
*,
entity_type: str,
source_record_id: str | None,
error: Exception,
raw_payload: dict | None = None,
) -> None:
"""抓取失败时尽力写入死信表(失败不影响主流程)。
死信是「可观测性」基础设施,与 RawEvent/Lineage 同级:写入失败只记
warning,绝不能让原始抓取错误之外的新异常打断采集循环。
"""
try:
await _write_ingest_failure(
db, "bzzoiro", entity_type, source_record_id,
"fetch_error", str(error), raw_payload,
)
except Exception:
logger.warning(
"写入 ingest_failures 死信失败(entity=%s, record=%s): %s",
entity_type, source_record_id, error, exc_info=True,
)
async def _write_lineage(db, source_system: str, source_record_id: str, target_table: str, target_id: int | None, transform_name: str, transform_detail: dict | None = None, batch_id: str | None = None) -> None: async def _write_lineage(db, source_system: str, source_record_id: str, target_table: str, target_id: int | None, transform_name: str, transform_detail: dict | None = None, batch_id: str | None = None) -> None:
"""写入 ETL 血缘追踪。""" """写入 ETL 血缘追踪。"""
db.add(DataLineage( db.add(DataLineage(
@@ -380,6 +439,53 @@ async def _write_lineage(db, source_system: str, source_record_id: str, target_t
)) ))
def _events_record_id(league_code: str, nm, raw: dict) -> str:
"""events 载荷的 RawEvent 幂等键。
优先用上游 event id;缺失时用 (league:home:away:date) 合成稳定键 ——
取 normalize 后的队名与天级日期(与 _match_key 同口径),不依赖 DB 自增 id,
保证同一来源比赛重复采集时命中同一条 RawEvent,不产生重复原始载荷。
"""
eid = _to_int_or_none(raw.get("id"))
if eid is not None:
return str(eid)
d = _to_date(nm.date)
date_part = d.isoformat() if d is not None else "na"
return f"{league_code}:{nm.home_team}:{nm.away_team}:{date_part}"
async def _write_events_bronze(
db,
*,
source_record_id: str,
raw_payload: dict,
target_match_id: int | None,
league_code: str,
match_status: str | None,
batch_id: str,
) -> None:
"""events 成功插入/更新单场比赛后的 Bronze 层补写:RawEvent(幂等) + DataLineage。
D1(工程债):此前只有 stats 回填写 RawEvent/Lineage,events 管线作为比赛
主数据的唯一入口反而不留溯源记录。幂等性由 _write_raw_event 的
source_record_id 查重保证;best-effort:基础设施写入失败只记 warning,
绝不拖垮采集主流程(与 _safe_write_ingest_failure 同级约束)。
"""
try:
await _write_raw_event(db, "bzzoiro", source_record_id, raw_payload, batch_id)
await _write_lineage(
db, "bzzoiro", source_record_id,
"matches", target_match_id, "events_ingest",
{"league": league_code, "match_status": match_status},
batch_id,
)
except Exception:
logger.warning(
"events Bronze 写入失败(record=%s, match=%s),不影响采集主流程",
source_record_id, target_match_id, exc_info=True,
)
# ============================================================ # ============================================================
# 积分榜管线:/leagues/{id}/standings/ → standings 表 # 积分榜管线:/leagues/{id}/standings/ → standings 表
# ============================================================ # ============================================================
@@ -418,12 +524,19 @@ async def ingest_bzzoiro_standings(db, *, leagues: Iterable[str], season: str |
result: dict = {"leagues": {}, "total_upserted": 0, "errors": []} result: dict = {"leagues": {}, "total_upserted": 0, "errors": []}
for code in leagues: for code in leagues:
league_r: dict = {"upserted": 0, "teams_created": 0, "rows": 0} league_r: dict = {"upserted": 0, "teams_created": 0, "rows": 0, "errors": []}
try: try:
payload = await fetch_bzzoiro_standings(code, season=season) payload = await fetch_bzzoiro_standings(code, season=season)
except Exception as e: except Exception as e:
logger.exception("bzzoiro standings fetch failed for %s", code) logger.exception("bzzoiro standings fetch failed for %s", code)
league_r["errors"].append(str(e)) league_r["errors"].append(str(e))
await _safe_write_ingest_failure(
db,
entity_type="standings",
source_record_id=None,
error=e,
raw_payload={"league": code, "season": season},
)
result["leagues"][code] = league_r result["leagues"][code] = league_r
result["errors"].append(f"{code}: {e}") result["errors"].append(f"{code}: {e}")
continue continue
@@ -434,13 +547,11 @@ async def ingest_bzzoiro_standings(db, *, leagues: Iterable[str], season: str |
result["errors"].append(f"{code}: 无积分榜数据") result["errors"].append(f"{code}: 无积分榜数据")
continue continue
# 联赛(get-or-create) # 联赛(get-or-create,D4: 经 LeagueRepository)
stmt = select(League).where(League.code == code) league = await LeagueRepository(db).get_or_create(
league = (await db.execute(stmt)).scalar_one_or_none() code, LEAGUE_NAMES.get(code, code), LEAGUE_COUNTRIES.get(code)
if league is None: )
league = League(code=code, name=LEAGUE_NAMES.get(code, code), country=LEAGUE_COUNTRIES.get(code)) team_r = TeamRepository(db)
db.add(league)
await db.flush()
# 赛季标签:优先用返回的 season 对象推导 # 赛季标签:优先用返回的 season 对象推导
season_obj = payload.get("season") or {} season_obj = payload.get("season") or {}
@@ -453,11 +564,7 @@ async def ingest_bzzoiro_standings(db, *, leagues: Iterable[str], season: str |
# 批量预载球队(与 events 管线使用同一 normalize 规则,保证 Team 匹配) # 批量预载球队(与 events 管线使用同一 normalize 规则,保证 Team 匹配)
names = {normalize_name(str(r.get("team_name", ""))) for r in rows} names = {normalize_name(str(r.get("team_name", ""))) for r in rows}
names.discard("") names.discard("")
team_map: dict[str, Team] = {} team_map: dict[str, Team] = await team_r.get_all_by_names(list(names))
if names:
stmt = select(Team).where(Team.name.in_(names))
for t in (await db.execute(stmt)).scalars():
team_map[t.name] = t
now = datetime.now(timezone.utc) now = datetime.now(timezone.utc)
for r in rows: for r in rows:
@@ -466,9 +573,7 @@ async def ingest_bzzoiro_standings(db, *, leagues: Iterable[str], season: str |
continue continue
team = team_map.get(team_name) team = team_map.get(team_name)
if team is None: if team is None:
team = Team(name=team_name, name_zh=zh_name(team_name)) team = await team_r.get_or_create(team_name, name_zh=zh_name(team_name))
db.add(team)
await db.flush()
team_map[team_name] = team team_map[team_name] = team
league_r["teams_created"] += 1 league_r["teams_created"] += 1
@@ -609,16 +714,10 @@ async def ingest_bzzoiro_event_stats(
result["errors"].append("无有效联赛代码") result["errors"].append("无有效联赛代码")
return result return result
stmt = ( # D4: 候选比赛查询经 MatchRepository(含 stats 预加载,筛选/排序/limit 语义不变)
select(Match) matches = await MatchRepository(db).find_finished_with_stats(
.options(selectinload(Match.stats)) league_ids, limit=limit * 3 if only_missing else limit
.where(Match.match_status == "finished")
.where(Match.source_event_id.is_not(None))
.where(Match.league_id.in_(league_ids))
.order_by(Match.match_date.desc())
.limit(limit * 3 if only_missing else limit)
) )
matches = (await db.execute(stmt)).scalars().all()
now = datetime.now(timezone.utc) now = datetime.now(timezone.utc)
processed = 0 processed = 0
@@ -634,6 +733,13 @@ async def ingest_bzzoiro_event_stats(
except Exception as e: except Exception as e:
logger.warning("stats fetch failed match=%s event=%s: %s", m.id, m.source_event_id, e) logger.warning("stats fetch failed match=%s event=%s: %s", m.id, m.source_event_id, e)
result["errors"].append(f"match {m.id}: {e}") result["errors"].append(f"match {m.id}: {e}")
await _safe_write_ingest_failure(
db,
entity_type="match_stats",
source_record_id=str(m.source_event_id),
error=e,
raw_payload={"match_id": m.id},
)
await asyncio.sleep(REQUEST_INTERVAL) await asyncio.sleep(REQUEST_INTERVAL)
continue continue
+2 -1
View File
@@ -1,7 +1,8 @@
"""API Key 轮换环:多 key 自动切换,遇到限流(429)自动跳过已冷却 key。 """API Key 轮换环:多 key 自动切换,遇到限流(429)自动跳过已冷却 key。
设计: 设计:
- 进程内纯内存状态(限速是短时状态,无需持久化) - 进程内纯内存状态(限速是短时状态,无需持久化;D7: 多 worker 部署时各进程
独立计数、不共享,上游限速额度应按 worker 数分摊,或前置网关统一管理)
- 单 key 场景零开销:直接透传 - 单 key 场景零开销:直接透传
- 多 key 场景:429 时把当前 key 标记冷却(默认 60s),轮转到下一个可用 key - 多 key 场景:429 时把当前 key 标记冷却(默认 60s),轮转到下一个可用 key
- 全部 key 都在冷却时:使用最早冷却的那个 key 并等待(退化到单 key 重试) - 全部 key 都在冷却时:使用最早冷却的那个 key 并等待(退化到单 key 重试)
+7 -3
View File
@@ -232,7 +232,9 @@ class Prediction(Base):
raw_response: Mapped[dict | None] = mapped_column(JSONB) raw_response: Mapped[dict | None] = mapped_column(JSONB)
# multi-agent 模式: 各专家报告 # multi-agent 模式: 各专家报告
mode: Mapped[str] = mapped_column(String(20), nullable=False, default="single") mode: Mapped[str] = mapped_column(String(20), nullable=False, default="single")
agent_outputs: Mapped[dict | None] = mapped_column(JSONB) # D5: multi-agent 模式存各专家报告列表(list[dict]);历史数据/兼容路径可能存 dict。
# 仅修正类型标注与真实 JSON 形状一致,列类型(JSONB)与数据不变。
agent_outputs: Mapped[list[dict] | dict | None] = mapped_column(JSONB)
# Fix: agent_weights 独立持久化到列(原本只在 raw_response 中) # Fix: agent_weights 独立持久化到列(原本只在 raw_response 中)
agent_weights: Mapped[dict | None] = mapped_column(JSONB) agent_weights: Mapped[dict | None] = mapped_column(JSONB)
# 预测状态: success / failed / degraded # 预测状态: success / failed / degraded
@@ -325,8 +327,10 @@ class RawEvent(Base):
class IngestFailure(Base): class IngestFailure(Base):
"""采集失败死信:记录失败原因、重试次数与下次重试时间。 """采集失败死信:记录失败原因、重试次数与下次重试时间。
⚠️ 预留未启用:当前 bzzoiro 管线不写入此表。 bzzoiro 三条管线(events / standings / stats)抓取失败时经由
未来接线计划:采集失败时写入,支持按 next_retry_at 自动重试。 bzzoiro._safe_write_ingest_failure 写入本表(尽力而为,写入失败
不影响采集主流程)。admin 可通过 /admin/schedules/ingest-failures
查看与重试。
""" """
__tablename__ = "ingest_failures" __tablename__ = "ingest_failures"
+31 -3
View File
@@ -64,6 +64,34 @@ class MatchRepository:
) )
return (await self._session.execute(stmt)).scalar_one_or_none() return (await self._session.execute(stmt)).scalar_one_or_none()
async def find_by_league_and_date_range(
self, league_id: int, start, end
) -> list[Match]:
"""批量预加载某联赛日期范围内的比赛(ingest 管线内存去重用)。"""
stmt = (
select(Match)
.where(Match.league_id == league_id)
.where(Match.match_date >= start)
.where(Match.match_date <= end)
)
return (await self._session.execute(stmt)).scalars().all()
async def find_finished_with_stats(self, league_ids: list[int], *, limit: int) -> list[Match]:
"""已完赛且有上游 event id 的比赛(按日期倒序),供统计回填逐场拉取。
预加载 stats:调用方需读取 existing.stats 判断是否跳过。
"""
stmt = (
select(Match)
.options(selectinload(Match.stats))
.where(Match.match_status == "finished")
.where(Match.source_event_id.is_not(None))
.where(Match.league_id.in_(league_ids))
.order_by(Match.match_date.desc())
.limit(limit)
)
return (await self._session.execute(stmt)).scalars().all()
async def add(self, match: Match) -> None: async def add(self, match: Match) -> None:
self._session.add(match) self._session.add(match)
await self._session.flush() await self._session.flush()
@@ -79,11 +107,11 @@ class TeamRepository:
stmt = select(Team).where(Team.name == name) stmt = select(Team).where(Team.name == name)
return (await self._session.execute(stmt)).scalar_one_or_none() return (await self._session.execute(stmt)).scalar_one_or_none()
async def get_or_create(self, name: str) -> Team: async def get_or_create(self, name: str, *, name_zh: str | None = None) -> Team:
"""按名获取球队,不存在则创建。""" """按名获取球队,不存在则创建(name_zh 供 bzzoiro 管线写中文名)"""
team = await self.get_by_name(name) team = await self.get_by_name(name)
if team is None: if team is None:
team = Team(name=name) team = Team(name=name, name_zh=name_zh)
self._session.add(team) self._session.add(team)
await self._session.flush() await self._session.flush()
return team return team
+7 -24
View File
@@ -6,14 +6,13 @@ import hashlib
import json import json
import logging import logging
import time import time
from dataclasses import dataclass
from datetime import datetime, timezone from datetime import datetime, timezone
from src.core.config import settings from src.core.config import settings
from src.db.base import AsyncSessionLocal from src.db.base import AsyncSessionLocal
from src.db.models import Match, Prediction from src.db.models import Match, Prediction
from src.db.unit_of_work import get_uow from src.db.unit_of_work import get_uow
from src.llm.predict import _upsert_prediction from src.llm.predict import PredictResult, _upsert_prediction
from src.llm.agents.base import AgentReport, AgentSpec, load_agent_prompt from src.llm.agents.base import AgentReport, AgentSpec, load_agent_prompt
from src.llm.context_builder import ( from src.llm.context_builder import (
MatchHeader, MatchHeader,
@@ -81,28 +80,12 @@ AGENT_LABELS_ZH: dict[str, str] = {
} }
@dataclass # D2(工程债): multi 结果类型与 single 统一 —— 扩展后的 PredictResult 用可选
class MultiPredictResult: # 字段(agent_outputs/agent_weights/prompt_tokens/completion_tokens/mode)承载
prediction_id: int # 全部模式,此处仅保留别名。保留 `MultiPredictResult` 名字的原因:
provider: str # 1. predict_match_multi 签名 `-> MultiPredictResult:` 是 R4 源码守卫的标记;
model: str # 2. src/llm/agents/__init__.py 对外 re-export 该名字。
prompt_version: str MultiPredictResult = PredictResult
mode: str
pred_home_goals: float | None
pred_away_goals: float | None
alt_pred_home_goals: int | None
alt_pred_away_goals: int | None
pred_1x2: str | None
subjective_confidence: float | None
reasoning: str | None
context: str
agent_outputs: list[dict]
agent_weights: dict | None
status: str = "success"
latency_ms: int | None = None
prompt_tokens: int | None = None
completion_tokens: int | None = None
raw: dict | None = None
async def _agent_provider(agent_id: str, *, tier: str, model_override: str | None = None) -> LLMProvider: async def _agent_provider(agent_id: str, *, tier: str, model_override: str | None = None) -> LLMProvider:
+25 -20
View File
@@ -12,6 +12,7 @@ from sqlalchemy import case, func, select
from src.db.base import AsyncSession, AsyncSessionLocal from src.db.base import AsyncSession, AsyncSessionLocal
from src.db.models import Match from src.db.models import Match
from src.llm.predict import PredictResult
logger = logging.getLogger(__name__) logger = logging.getLogger(__name__)
@@ -52,11 +53,13 @@ async def predict_baseline(
*, *,
backtest: bool = False, backtest: bool = False,
cutoff_at: datetime | None = None, cutoff_at: datetime | None = None,
) -> dict: ) -> PredictResult:
"""极简基线预测:主场场均进球 vs 客场场均进球。 """极简基线预测:主场场均进球 vs 客场场均进球。
返回 PredictResult 兼容的字典: 返回 PredictResult(D2 统一结果类型):
provider=model="baseline", 不调用 LLM,latency_ms≈0。 provider=model="baseline", 不调用 LLM,latency_ms≈0。
prediction_id 为占位 0 —— baseline 不在服务层落库,
由路由层 _persist_baseline 落库后取得真实 id。
""" """
async with AsyncSessionLocal() as db: async with AsyncSessionLocal() as db:
match = await db.get(Match, match_id) match = await db.get(Match, match_id)
@@ -90,24 +93,26 @@ async def predict_baseline(
else: else:
pred_1x2 = "X" pred_1x2 = "X"
return { return PredictResult(
"pred_home_goals": float(pred_home), prediction_id=0, # 占位:真实 id 由路由层 _persist_baseline 落库后返回
"pred_away_goals": float(pred_away), provider="baseline",
"alt_pred_home_goals": None, model="baseline",
"alt_pred_away_goals": None, prompt_version="baseline_v1",
"pred_1x2": pred_1x2, mode="baseline",
"subjective_confidence": 0.5, pred_home_goals=float(pred_home),
"prompt_tokens": 0, pred_away_goals=float(pred_away),
"completion_tokens": 0, alt_pred_home_goals=None,
"reasoning": ( alt_pred_away_goals=None,
pred_1x2=pred_1x2,
subjective_confidence=0.5,
reasoning=(
f"基线估计(非投注建议): 主队主场场均进球 {home_avg:.2f} → 预测 {pred_home}; " f"基线估计(非投注建议): 主队主场场均进球 {home_avg:.2f} → 预测 {pred_home}; "
f"客队客场场均进球 {away_avg:.2f} → 预测 {pred_away}" f"客队客场场均进球 {away_avg:.2f} → 预测 {pred_away}"
), ),
"provider": "baseline", context="", # baseline 不构建 LLM 上下文
"model": "baseline", status="success",
"prompt_version": "baseline_v1", latency_ms=0,
"mode": "baseline", prompt_tokens=0,
"status": "success", completion_tokens=0,
"latency_ms": 0, raw={"home_avg": round(home_avg, 2), "away_avg": round(away_avg, 2)},
"raw": {"home_avg": round(home_avg, 2), "away_avg": round(away_avg, 2)}, )
}
+18 -1
View File
@@ -89,6 +89,14 @@ def _prompt_template_hash(version: str) -> str:
@dataclass @dataclass
class PredictResult: class PredictResult:
"""三种预测模式(single/multi/baseline)的统一结果类型。
D2(工程债): 原本 single 返回本类、multi 重复定义 MultiPredictResult、
baseline 返回裸 dict,导致路由 isinstance(dict) 双分支 + backtest 对
baseline 直接 AttributeError。现以可选字段扩展本类承载全部模式;
MultiPredictResult 是本类的别名(见 src/llm/agents/orchestrator.py)。
"""
prediction_id: int prediction_id: int
provider: str provider: str
model: str model: str
@@ -101,8 +109,15 @@ class PredictResult:
subjective_confidence: float | None subjective_confidence: float | None
reasoning: str | None reasoning: str | None
context: str context: str
# 模式标识: single(默认) / multi / baseline
mode: str = "single"
# multi 专属: 各专家报告列表与融合权重;single/baseline 为 None
agent_outputs: list[dict] | None = None
agent_weights: dict | None = None
status: str = "success" status: str = "success"
latency_ms: int | None = None latency_ms: int | None = None
prompt_tokens: int | None = None
completion_tokens: int | None = None
raw: dict | None = None raw: dict | None = None
@@ -158,9 +173,11 @@ async def predict_match(
use_cache: bool = True, use_cache: bool = True,
backtest: bool = False, backtest: bool = False,
cutoff_at=None, cutoff_at=None,
) -> "PredictResult | MultiPredictResult": ) -> PredictResult:
"""预测入口。mode=multi(默认)走多 agent;mode=single 走单次调用;mode=baseline 走无 LLM 基线。 """预测入口。mode=multi(默认)走多 agent;mode=single 走单次调用;mode=baseline 走无 LLM 基线。
三种模式统一返回 PredictResult(D2);multi 的 MultiPredictResult 是其别名。
Args: Args:
mode: multi(默认,5 专家+终裁) / single(单次) / baseline(极简统计基线,不调用 LLM)。 mode: multi(默认,5 专家+终裁) / single(单次) / baseline(极简统计基线,不调用 LLM)。
use_cache:是否允许返回进程内缓存结果。回测必须传 False—— use_cache:是否允许返回进程内缓存结果。回测必须传 False——
+14 -14
View File
@@ -81,18 +81,18 @@ async def test_predict_baseline_no_llm():
result = await predict_baseline(1) result = await predict_baseline(1)
assert result["provider"] == "baseline" assert result.provider == "baseline"
assert result["model"] == "baseline" assert result.model == "baseline"
assert result["mode"] == "baseline" assert result.mode == "baseline"
assert result["latency_ms"] == 0 assert result.latency_ms == 0
assert result["prompt_tokens"] == 0 assert result.prompt_tokens == 0
assert result["completion_tokens"] == 0 assert result.completion_tokens == 0
# 2.4 → round = 2, 1.6 → round = 2 → 平局 X # 2.4 → round = 2, 1.6 → round = 2 → 平局 X
assert result["pred_home_goals"] == 2.0 assert result.pred_home_goals == 2.0
assert result["pred_away_goals"] == 2.0 assert result.pred_away_goals == 2.0
assert result["pred_1x2"] == "X" assert result.pred_1x2 == "X"
assert result["subjective_confidence"] == 0.5 assert result.subjective_confidence == 0.5
assert "非投注建议" in result["reasoning"] assert "非投注建议" in result.reasoning
# 确认未调用任何 LLM 相关模块 # 确认未调用任何 LLM 相关模块
assert "home_10" in captured and "away_20" in captured assert "home_10" in captured and "away_20" in captured
@@ -125,6 +125,6 @@ async def test_predict_baseline_clamps_to_range():
result = await predict_baseline(2) result = await predict_baseline(2)
assert result["pred_home_goals"] == 10.0 # clamped assert result.pred_home_goals == 10.0 # clamped
assert result["pred_away_goals"] == 0.0 # clamped assert result.pred_away_goals == 0.0 # clamped
assert result["pred_1x2"] == "1" # 10:0 主胜 assert result.pred_1x2 == "1" # 10:0 主胜
+216
View File
@@ -0,0 +1,216 @@
"""D2 工程债回归测试: 统一预测结果类型。
背景: predict_match 三条路径返回类型不一 —— single 返回 PredictResult,
multi 返回字段重复定义的 MultiPredictResult dataclass,baseline 返回裸 dict。
后果: (1) 预测路由 PredictOut 映射被迫写 isinstance(result, dict) 双分支;
(2) backtest 对 baseline 模式直接 AttributeError(dict 没有 .prediction_id,
潜伏 bug);(3) 字段清单在两处 dataclass 重复维护,加字段必漏一处。
统一方案: 扩展 PredictResult(可选字段)承载全部模式;
MultiPredictResult 变为其别名(保留 orchestrator 签名标记,兼容 re-export);
baseline 返回 PredictResult;路由单一字段映射。
本测试守护四件事:
1. predict_baseline 返回 PredictResult(属性访问)
2. MultiPredictResult 与 PredictResult 兼容(orchestrator 构造调用的
全字段 kwargs 可直接构造别名)
3. 预测路由不再有 isinstance(result, dict) 分支(源码守卫,仿 R4 范式)
4. _persist_baseline 用属性访问构造 upsert values(baseline 落库语义不变)
"""
from __future__ import annotations
from pathlib import Path
from types import SimpleNamespace
import pytest
from src.llm.baseline import predict_baseline
from src.llm.predict import PredictResult
from src.llm.agents import MultiPredictResult
ROUTE_PATH = Path(__file__).resolve().parents[1] / "src" / "api" / "routes" / "predict.py"
# ============================================================
# 1. baseline 返回 PredictResult
# ============================================================
@pytest.mark.asyncio
async def test_predict_baseline_returns_predict_result():
"""基线预测返回 PredictResult 实例,mode=baseline,token/延迟为 0。"""
from unittest.mock import patch
async def fake_avg(db, *, team_id, side, league_id, before):
return 2.4 if side == "home" else 1.6
class FakeMatch:
id = 1
home_team_id = 10
away_team_id = 20
league_id = 1
match_status = "scheduled"
class FakeSession:
async def get(self, cls, mid):
return FakeMatch()
class FakeCM:
async def __aenter__(self):
return FakeSession()
async def __aexit__(self, *a):
return None
with patch("src.llm.baseline._avg_goals", fake_avg), \
patch("src.llm.baseline.AsyncSessionLocal") as SLC:
SLC.return_value = FakeCM()
result = await predict_baseline(1)
assert isinstance(result, PredictResult)
assert result.mode == "baseline"
assert result.provider == "baseline"
assert result.model == "baseline"
assert result.prompt_version == "baseline_v1"
assert result.pred_home_goals == 2.0
assert result.pred_away_goals == 2.0
assert result.pred_1x2 == "X"
assert result.subjective_confidence == 0.5
assert result.prompt_tokens == 0
assert result.completion_tokens == 0
assert result.latency_ms == 0
assert result.status == "success"
assert "非投注建议" in (result.reasoning or "")
# baseline 不构建 LLM 上下文,但字段必须存在且可安全序列化
assert result.context == ""
# 原始统计快照保留在 raw 中
assert result.raw is not None
assert "home_avg" in result.raw
# ============================================================
# 2. MultiPredictResult 与扩展后的 PredictResult 兼容
# ============================================================
def test_multi_predict_result_is_predict_result_alias():
"""multi 结果不再是重复定义的 dataclass,而是扩展 PredictResult 的别名。"""
assert MultiPredictResult is PredictResult
def test_multi_result_constructor_kwargs_still_supported():
"""orchestrator 现有构造调用的全部字段 kwargs 必须仍可构造(别名完整性)。"""
# 与 orchestrator.predict_match_multi 的 return MultiPredictResult(...) 逐一对应
result = MultiPredictResult(
prediction_id=1,
provider="openai",
model="gpt-x",
prompt_version="multi_v1",
mode="multi",
pred_home_goals=2.0,
pred_away_goals=1.0,
alt_pred_home_goals=None,
alt_pred_away_goals=None,
pred_1x2="1",
subjective_confidence=0.7,
reasoning="r",
status="success",
agent_outputs=[{"agent": "form"}],
agent_weights={"form": 0.2},
context="ctx",
latency_ms=100,
prompt_tokens=10,
completion_tokens=5,
raw={"final": True},
)
assert result.mode == "multi"
assert result.agent_outputs == [{"agent": "form"}]
assert result.agent_weights == {"form": 0.2}
assert result.prompt_tokens == 10
assert result.completion_tokens == 5
# ============================================================
# 3. 路由去 dict 分支(源码守卫)
# ============================================================
def test_predict_route_has_no_dict_branch():
"""PredictOut 映射必须统一走属性访问,禁止 isinstance(result, dict) 回潮。"""
src = ROUTE_PATH.read_text(encoding="utf-8")
assert "isinstance(result, dict)" not in src
assert ".get(\"pred_home_goals\")" not in src
# ============================================================
# 4. _persist_baseline 属性映射(baseline 落库语义不变)
# ============================================================
class _FakeUoW:
"""替代 get_uow 的最小上下文管理器。"""
def __init__(self):
self.session = SimpleNamespace()
async def __aenter__(self):
return self.session
async def __aexit__(self, *a):
return None
@pytest.mark.asyncio
async def test_persist_baseline_maps_attributes(monkeypatch):
captured = {}
async def fake_upsert(session, **kwargs):
captured.update(kwargs)
return SimpleNamespace(id=77)
monkeypatch.setattr("src.db.unit_of_work.get_uow", lambda: _FakeUoW())
monkeypatch.setattr("src.llm.predict._upsert_prediction", fake_upsert)
from src.api.routes.predict import _persist_baseline
baseline = PredictResult(
prediction_id=0, # baseline 不在服务层落库,由 _persist_baseline 落库后取得真实 id
provider="baseline",
model="baseline",
prompt_version="baseline_v1",
mode="baseline",
pred_home_goals=2.0,
pred_away_goals=1.0,
alt_pred_home_goals=None,
alt_pred_away_goals=None,
pred_1x2="1",
subjective_confidence=0.5,
reasoning="r",
context="",
status="success",
latency_ms=0,
prompt_tokens=0,
completion_tokens=0,
raw={"home_avg": 2.1, "away_avg": 1.4},
)
pid = await _persist_baseline(1, baseline)
assert pid == 77
assert captured["match_id"] == 1
assert captured["provider_name"] == "baseline"
assert captured["model"] == "baseline"
assert captured["mode"] == "baseline"
assert captured["run_type"] == "live"
v = captured["values"]
assert v["prompt_version"] == "baseline_v1"
assert v["pred_home_goals"] == 2.0
assert v["pred_away_goals"] == 1.0
assert v["pred_1x2"] == "1"
assert v["subjective_confidence"] == 0.5
assert v["prompt_tokens"] == 0
assert v["completion_tokens"] == 0
assert v["latency_ms"] == 0
assert v["raw_response"] == {"home_avg": 2.1, "away_avg": 1.4}
assert v["status"] == "success"
+259
View File
@@ -0,0 +1,259 @@
"""D1 工程债回归测试: events 成功路径必须写 Bronze 层(RawEvent + DataLineage)。
背景: stats 回填管线早已有 RawEvent/DataLineage 写入,但 events 管线(比赛
主数据的唯一入口)成功插入/更新后既不留原始载荷,也不留血缘 —— 数据溯源
链条在最关键的一环断掉。本测试守护:
1. 插入新比赛 → RawEvent(幂等键=source_event_id 或合成键) + Lineage
(target_table="matches", transform_name="events_ingest")
2. 变更更新(如补比分/状态) → 同样写血缘
3. 无变化跳过 → 不写(避免 lineage 刷屏)
4. RawEvent 幂等: 同 source_record_id 已存在则跳过
5. 基础设施写入失败 → 只 warning,不拖垮采集主流程
范式: 假 db(按查询实体分发预置数据 + 记录 add,flush 分配自增 id)
+ monkeypatch 抓取函数,不依赖真实数据库。
"""
from __future__ import annotations
from datetime import date, datetime, timezone
import pytest
import src.data.bzzoiro as bz
from src.db.models import DataLineage, League, Match, RawEvent, Team
def _event(eid=1001, status="finished", home="Arsenal", away="Chelsea", hs=2, as_=1):
"""构造一条最小合法的 bzzoiro /events/ 原始载荷。"""
raw = {
"event_date": "2026-09-20 15:00:00",
"status": status,
"home_team": home,
"away_team": away,
"home_score": hs,
"away_score": as_,
}
if eid is not None:
raw["id"] = eid
return raw
class _FakeResult:
"""支持 .scalars().all() / .scalar_one_or_none() 的最小假结果集。"""
def __init__(self, items):
self._items = list(items)
def scalars(self):
return self
def __iter__(self):
return iter(self._items)
def all(self):
return self._items
def scalar_one_or_none(self):
return self._items[0] if self._items else None
def scalar(self):
return None
class _FakeDB:
"""按查询实体分发预置数据;记录 add();flush 为无 id 对象分配自增主键。"""
def __init__(self, matches=(), teams=(), leagues=(), raw_events=()):
self.added = []
self._by_entity = {
Match: list(matches),
Team: list(teams),
League: list(leagues),
RawEvent: list(raw_events),
}
self._next_id = 0
def add(self, obj):
self.added.append(obj)
async def execute(self, stmt):
entities = set()
for d in (stmt.column_descriptions or []):
entities.add(d.get("entity") or d.get("type"))
for entity, items in self._by_entity.items():
if entity in entities:
return _FakeResult(items)
return _FakeResult([])
async def flush(self):
for obj in self.added:
if getattr(obj, "id", None) is None:
self._next_id += 1
obj.id = self._next_id
@pytest.fixture(autouse=True)
def _no_request_interval(monkeypatch):
monkeypatch.setattr(bz, "REQUEST_INTERVAL", 0)
def _patch_fetch(monkeypatch, events):
async def _fetch(league_code, **kwargs):
return list(events)
monkeypatch.setattr(bz, "fetch_bzzoiro_events", _fetch)
def _matches(db):
return [o for o in db.added if isinstance(o, Match)]
def _raw_events(db):
return [o for o in db.added if isinstance(o, RawEvent)]
def _lineages(db):
return [o for o in db.added if isinstance(o, DataLineage)]
# ============================================================
# 1. 插入新比赛 → RawEvent + DataLineage
# ============================================================
class TestEventsBronzeOnInsert:
async def test_insert_writes_raw_event_and_lineage(self, monkeypatch):
_patch_fetch(monkeypatch, [_event()])
db = _FakeDB()
result = await bz.BzzoiroSource().ingest(db, leagues=["E0"])
assert result["total_inserted"] == 1
raws = _raw_events(db)
assert len(raws) == 1
raw = raws[0]
assert raw.source_system == "bzzoiro"
assert raw.source_record_id == "1001" # 有上游 id 时直接用
assert raw.ingest_batch_id.startswith("bzzoiro-events-E0-")
assert raw.raw_payload["id"] == 1001 # 原始载荷完整保留
lineages = _lineages(db)
assert len(lineages) == 1
lin = lineages[0]
assert lin.source_system == "bzzoiro"
assert lin.source_record_id == "1001"
assert lin.target_table == "matches"
assert lin.target_id == _matches(db)[0].id
assert lin.transform_name == "events_ingest"
# RawEvent 与 Lineage 同批次,便于按批追溯
assert lin.batch_id == raw.ingest_batch_id
async def test_missing_source_id_uses_synthetic_stable_key(self, monkeypatch):
"""上游 id 缺失时,用 (league:home:away:date) 合成稳定幂等键。"""
_patch_fetch(monkeypatch, [_event(eid=None)])
db = _FakeDB()
result = await bz.BzzoiroSource().ingest(db, leagues=["E0"])
assert result["total_inserted"] == 1
raws = _raw_events(db)
assert len(raws) == 1
# 期望键基于 normalize 后的队名与天级日期 —— 与 _match_key 同口径,
# 不依赖 DB 自增 id,跨批次可复现
nm = bz.normalize_bzzoiro(_event(eid=None), "E0")
expected = f"E0:{nm.home_team}:{nm.away_team}:{nm.date.date().isoformat()}"
assert raws[0].source_record_id == expected
async def test_existing_raw_event_is_skipped(self, monkeypatch):
"""RawEvent 幂等: 同 source_record_id 已存在则不再新增,但血缘照写。"""
existing = RawEvent(
source_system="bzzoiro",
source_record_id="1001",
raw_payload={"old": True},
)
_patch_fetch(monkeypatch, [_event()])
db = _FakeDB(raw_events=[existing])
await bz.BzzoiroSource().ingest(db, leagues=["E0"])
new_raws = [r for r in _raw_events(db) if r is not existing]
assert new_raws == []
assert len(_lineages(db)) == 1 # 血缘仍然记录本次采集
# ============================================================
# 2. 变更更新 → 写血缘;无变化 → 不写
# ============================================================
class TestEventsBronzeOnUpdate:
def _existing_match(self, **overrides):
m = Match(
league_id=1,
home_team_id=2, # 与本轮 Team 创建后 fake 自增 id 对齐(league=1, home=2, away=3)
away_team_id=3,
match_date=datetime(2026, 9, 20, 15, 0, tzinfo=timezone.utc),
match_date_date=date(2026, 9, 20),
match_status="scheduled",
source_event_id=1001,
)
m.id = 42
for k, v in overrides.items():
setattr(m, k, v)
return m
async def test_changed_update_writes_lineage(self, monkeypatch):
# 已有比赛处于 scheduled 且无比分;新载荷为 finished 2:1 → 触发变更更新
db = _FakeDB(matches=[self._existing_match()])
_patch_fetch(monkeypatch, [_event()])
result = await bz.BzzoiroSource().ingest(db, leagues=["E0"])
assert result["total_inserted"] == 0
assert result["leagues"]["E0"]["updated"] == 1
lineages = _lineages(db)
assert len(lineages) == 1
assert lineages[0].target_id == 42
assert lineages[0].target_table == "matches"
assert lineages[0].transform_name == "events_ingest"
async def test_unchanged_match_writes_nothing(self, monkeypatch):
# 已有比赛与新载荷完全一致 → 无变化,不应产生 RawEvent/Lineage
existing = self._existing_match(
match_status="finished",
home_goals=2,
away_goals=1,
)
db = _FakeDB(matches=[existing])
_patch_fetch(monkeypatch, [_event()])
result = await bz.BzzoiroSource().ingest(db, leagues=["E0"])
assert result["total_inserted"] == 0
assert result["leagues"]["E0"]["updated"] == 0
assert _raw_events(db) == []
assert _lineages(db) == []
# ============================================================
# 3. 基础设施写入失败: 尽力而为,不拖垮主流程
# ============================================================
class TestEventsBronzeIsBestEffort:
async def test_bronze_write_failure_does_not_break_ingest(self, monkeypatch):
async def _boom(*args, **kwargs):
raise RuntimeError("infra down")
monkeypatch.setattr(bz, "_write_raw_event", _boom)
monkeypatch.setattr(bz, "_write_lineage", _boom)
_patch_fetch(monkeypatch, [_event()])
db = _FakeDB()
# 不应抛异常:Bronze 写不进去只记 warning
result = await bz.BzzoiroSource().ingest(db, leagues=["E0"])
assert result["total_inserted"] == 1
assert len(_matches(db)) == 1
+209
View File
@@ -0,0 +1,209 @@
"""死信接线回归测试: bzzoiro 三条管线抓取失败必须写入 IngestFailure。
背景: IngestFailure 死信表此前「预留未启用」——events / standings / stats
抓取失败只打日志 + errors 列表,admin 的 /ingest-failures 列表与 retry
端点形同虚设。本测试守护三条管线的失败写入路径:
1. events 整联赛抓取失败 → entity_type="events"
2. standings 整联赛抓取失败 → entity_type="standings"
3. stats 单场统计抓取失败 → entity_type="match_stats"(带 source_record_id)
范式: 假 db(记录 add 调用) + monkeypatch 抓取函数,不依赖真实数据库 ——
与 test_review_required_fixes.py R2 相同。失败写入是「尽力而为」:
写入器自身抛错只记日志,不得拖垮采集主流程(最后一个测试守护)。
"""
from __future__ import annotations
from types import SimpleNamespace
import pytest
import src.data.bzzoiro as bz
from src.db.models import IngestFailure
class _FakeResult:
"""支持 .scalars().all() / .scalar_one_or_none() / .scalar() 的最小假结果集。"""
def __init__(self, items):
self._items = items
def scalars(self):
return self
def all(self):
return self._items
def scalar_one_or_none(self):
return None
def scalar(self):
return None
class _FakeDB:
"""只记录 add() 的假会话 —— 失败路径不触发真实查询。"""
def __init__(self, items=None):
self.added = []
self._items = items or []
def add(self, obj):
self.added.append(obj)
async def execute(self, stmt):
return _FakeResult(self._items)
async def flush(self):
pass
@pytest.fixture(autouse=True)
def _no_request_interval(monkeypatch):
"""失败路径会 await asyncio.sleep(REQUEST_INTERVAL),置 0 加速测试。"""
monkeypatch.setattr(bz, "REQUEST_INTERVAL", 0)
def _failures(db: _FakeDB) -> list[IngestFailure]:
return [o for o in db.added if isinstance(o, IngestFailure)]
# ============================================================
# 1. events: 整联赛抓取失败
# ============================================================
class TestEventsFetchFailureDeadLetter:
async def test_writes_deadletter_with_league_context(self, monkeypatch):
async def _boom(*args, **kwargs):
raise RuntimeError("network down")
monkeypatch.setattr(bz, "fetch_bzzoiro_events", _boom)
db = _FakeDB()
result = await bz.BzzoiroSource().ingest(db, leagues=["E0"])
rows = _failures(db)
assert len(rows) == 1
row = rows[0]
assert row.source_system == "bzzoiro"
assert row.entity_type == "events"
assert row.error_type == "fetch_error"
assert "network down" in row.error_detail
# 上下文足够管理员定位:联赛代码必须随行
assert row.raw_payload is not None
assert row.raw_payload.get("league") == "E0"
# 主流程不受影响:错误仍记录在 result 中
assert result["leagues"]["E0"]["errors"]
assert result["total_inserted"] == 0
async def test_other_leagues_continue_after_failure(self, monkeypatch):
async def _boom(league_code, **kwargs):
if league_code == "E0":
raise RuntimeError("boom")
return []
monkeypatch.setattr(bz, "fetch_bzzoiro_events", _boom)
db = _FakeDB()
await bz.BzzoiroSource().ingest(db, leagues=["E0", "SP1"])
rows = _failures(db)
assert len(rows) == 1
assert rows[0].raw_payload.get("league") == "E0"
# ============================================================
# 2. standings: 整联赛抓取失败
# ============================================================
class TestStandingsFetchFailureDeadLetter:
async def test_writes_deadletter_with_season_context(self, monkeypatch):
async def _boom(league_code, season=None):
raise RuntimeError("upstream 500")
monkeypatch.setattr(bz, "fetch_bzzoiro_standings", _boom)
db = _FakeDB()
result = await bz.ingest_bzzoiro_standings(db, leagues=["SP1"], season="2025-2026")
rows = _failures(db)
assert len(rows) == 1
row = rows[0]
assert row.entity_type == "standings"
assert row.error_type == "fetch_error"
assert "upstream 500" in row.error_detail
assert row.raw_payload == {"league": "SP1", "season": "2025-2026"}
assert result["errors"]
# ============================================================
# 3. stats: 单场统计抓取失败
# ============================================================
class TestStatsFetchFailureDeadLetter:
def _match(self) -> SimpleNamespace:
return SimpleNamespace(
id=42,
source_event_id=777,
stats=None,
match_date=None,
league_id=1,
match_status="finished",
)
async def test_writes_deadletter_with_source_record_id(self, monkeypatch):
async def _boom(path, params=None, max_retries=3):
raise TimeoutError("read timeout")
monkeypatch.setattr(bz, "_fetch_json_async", _boom)
db = _FakeDB(items=[self._match()])
result = await bz.ingest_bzzoiro_event_stats(db, leagues=["E0"], limit=1)
rows = _failures(db)
assert len(rows) == 1
row = rows[0]
assert row.entity_type == "match_stats"
# source_record_id 必须是上游 event id,retry 端点据此定位
assert row.source_record_id == "777"
assert "timeout" in (row.error_detail or "").lower()
assert row.raw_payload == {"match_id": 42}
assert result["errors"]
async def test_success_path_does_not_write_deadletter(self, monkeypatch):
async def _ok(path, params=None, max_retries=3):
return {"stats": {"home": {"total_shots": 10}, "away": {"total_shots": 5}}}
monkeypatch.setattr(bz, "_fetch_json_async", _ok)
db = _FakeDB(items=[self._match()])
await bz.ingest_bzzoiro_event_stats(db, leagues=["E0"], limit=1)
assert _failures(db) == []
# ============================================================
# 4. 死信写入自身失败: 尽力而为,不拖垮主流程
# ============================================================
class TestDeadLetterWriteIsBestEffort:
async def test_db_add_failure_does_not_break_ingest(self, monkeypatch):
async def _boom(*args, **kwargs):
raise RuntimeError("network down")
monkeypatch.setattr(bz, "fetch_bzzoiro_events", _boom)
class _BrokenDB(_FakeDB):
def add(self, obj):
if isinstance(obj, IngestFailure):
raise RuntimeError("session closed")
super().add(obj)
db = _BrokenDB()
# 不应抛异常:死信写不进去只记 warning
result = await bz.BzzoiroSource().ingest(db, leagues=["E0"])
assert result["leagues"]["E0"]["errors"]