Files
Profeto/alembic/versions/0019_ingest_jobs.py
WorkBuddy 983363dab7 feat(ingest): 采集任务状态跟踪(ingest_jobs),解决 fire-and-forget 不可观测
- 新表 ingest_jobs(迁移 0019): id UUID/task/params JSONB/
  status(pending|running|success|failed,CheckConstraint)/
  result JSONB(统计摘要)/error/created_at/started_at/finished_at
- POST /ingest/bzzoiro: 启动后台前创建 pending job,响应返回 job_id;
  仍 require_admin。后台 _run_bzzoiro 流转 running→success/failed,
  result 按子任务(events/standings/stats)记录摘要(errors 截断 10 条)
- _update_job 尽力而为: 状态更新失败只记日志,绝不拖垮采集主流程;
  与 IngestFailure 死信独立(行级 vs 任务级,可同时存在)
- 新增 admin 端点(挂 /api/v1/admin 路由,路由级 require_admin):
  GET /admin/ingest/jobs/{job_id} 与 GET /admin/ingest/jobs?limit&status
- 前端采集页: 提交后凭 job_id 3 秒轮询,终态展示结果摘要/失败原因;
  无 job_id 时回退旧的 30 秒盲等 + 系统日志提示
- 测试 11 项: 建 job+job_id 契约、非法 task 422、成功/失败/all 流转、
  update 失败不拖垮采集、部分失败仍 success、admin 端点 200/404/列表、
  结构守护(job 路由在 admin 路由且带 require_admin)
- 禁止项确认: 未动分批 UoW、BzzoiroSource、死信与 Bronze 写入;
  docs(01/03/05/07/README)同步 13 张表与端点说明
2026-09-21 21:19:21 +08:00

47 lines
1.6 KiB
Python

"""新增 ingest_jobs 表
Revision ID: 0019_ingest_jobs
Revises: 0018_match_checks
Create Date: 2026-09-21
采集任务状态跟踪:POST /ingest/bzzoiro 创建 pending job 后台执行,
解决 fire-and-forget 不可观测问题(admin 可查询任务级状态与结果摘要)。
"""
from typing import Sequence, Union
from alembic import op
import sqlalchemy as sa
from sqlalchemy.dialects.postgresql import JSONB
# revision identifiers, used by Alembic.
revision: str = '0019_ingest_jobs'
down_revision: Union[str, None] = '0018_match_checks'
branch_labels: Union[str, Sequence[str], None] = None
depends_on: Union[str, Sequence[str], None] = None
def upgrade() -> None:
op.create_table(
'ingest_jobs',
sa.Column('id', sa.String(36), primary_key=True),
sa.Column('task', sa.String(20), nullable=False),
sa.Column('params', JSONB, nullable=False),
sa.Column('status', sa.String(20), nullable=False, server_default='pending'),
sa.Column('result', JSONB),
sa.Column('error', sa.Text()),
sa.Column('created_at', sa.DateTime(timezone=True)),
sa.Column('started_at', sa.DateTime(timezone=True)),
sa.Column('finished_at', sa.DateTime(timezone=True)),
sa.CheckConstraint(
"status IN ('pending','running','success','failed')",
name='ck_ingest_jobs_status',
),
)
op.create_index('ix_ingest_jobs_created', 'ingest_jobs', ['created_at'])
def downgrade() -> None:
op.drop_index('ix_ingest_jobs_created', table_name='ingest_jobs')
op.drop_table('ingest_jobs')