DDevInfra 后端教程
Ctrl K
BACKEND · BEGINNER TO DEBUGGER · COMMIT e92d894

main.py 出发,亲手走完一个企业级 Agent 后端

这不是把名词堆给你的架构说明。你会跟着一条真实请求,依次看懂入口、API、Service、Repository、Worker、Agent、工具、证据、审批、RAG、评测与生产设施,最后能自己定位、调试和修改代码。

44 个 Python 源码文件 31 个 API 路由 144 个测试 离线模式先学 生产模式后拆
00

先学会怎么用这份教程

你的目标不是背代码,而是形成“看到现象 → 找到入口 → 沿调用链定位”的能力。

第一遍:只走主链

从第 01 章读到第 11 章。先接受“一个请求会经过很多层”,不要急着记每个方法。

第二遍:边跑边断

照第 03、15 章启动离线服务,在 create_runRunService.createRunWorker.execute 处打断点。

第三遍:改一个小功能

从修改配方挑一个任务,先找测试,再改实现。此时第 16 章是你的后端地图。

第一次看到陌生词怎么办?

先看它旁边的“新手补课”和变量表。页面顶部可搜索任何函数、变量、包或接口。代码里的 selfawaitDepends 等只在第一次出现时完整解释,后面会链接回来。

整条故事线

flowchart LR
  A["浏览器/Swagger"] --> B["FastAPI 路由"]
  B --> C["领域 Service"]
  C --> D["Repository"]
  C --> E["Dispatcher"]
  E --> F["Worker"]
  F --> G["Agent Runtime"]
  G --> H["Skill + Tool + Policy"]
  H --> I["Evidence / Approval"]
  F --> D
  D --> J["Trace / SSE / Result"]
            

图加载失败也不影响阅读:它表达的就是“请求进入 API,经服务和仓库存下任务,由 Worker 执行 Agent,最后把可审计结果写回仓库”。

01

项目目录与后端边界

先知道每个房间放什么,再进房间看家具。

devinfra-agent-platform/ ├─ backend/ ← 本教程主角:Python 后端 │ ├─ pyproject.toml ← 包信息、依赖、命令与质量工具 │ ├─ uv.lock ← 锁定可复现的精确依赖版本 │ ├─ Dockerfile ← 生产镜像怎么构建、以谁运行 │ ├─ alembic.ini / alembic/ ← 数据库版本迁移 │ ├─ scripts/ ← 导出 OpenAPI、负载测试 │ ├─ tests/ ← 144 个后端行为测试 │ └─ src/devinfra_agent/ ← 真正的业务源码包 │ ├─ main.py ← Uvicorn 导入的 API 入口 │ ├─ composition.py ← 离线/生产依赖装配 │ ├─ app.py ← FastAPI 路由与 AppServices │ ├─ services.py ← Run 与审批业务规则 │ ├─ repositories.py ← 内存/SQLAlchemy 数据访问 │ ├─ worker.py ← 后台任务执行与持久化投影 │ ├─ runtime.py / graph.py ← Agent 状态与 LangGraph │ ├─ skills.py / tools.py ← 工作流和工具目录 │ ├─ policy.py / evidence.py ← 权限决策和证据规范化 │ ├─ rag.py ← 切块、向量、BM25、RRF、Qdrant │ ├─ evaluation.py / harness/ ← 持久评测与离线评测场 │ └─ adapters/ ← MCP、GitHub、HTTP 外部边界 ├─ mcp/ ← TypeScript 工具服务;只讲 Python 如何调用 ├─ scenarios/ ← 固定案例、工具 fixture、72 道评测题 ├─ ops/ ← Prometheus、OTel、Grafana 等运维配置 ├─ frontend/ ← 本教程不展开其 React 实现 └─ compose.yaml ← 把全部生产服务接起来

三个容易混淆的词

项目(project)

整个仓库。里面既有 Python 后端,也有前端、MCP 服务、场景数据和运维文件。

包(package)

devinfra_agent 是可导入的 Python 包;目录里有 __init__.py

模块(module)

一个 .py 文件就是模块,例如 services.py;用 from .services import RunService 引入。

为什么使用 src/ 布局?

源码不直接躺在 backend/ 根目录,而在 backend/src/devinfra_agent。这样测试必须像真实安装后的用户一样导入包,可避免“只因当前目录碰巧在搜索路径里所以能运行”的假成功。

后端边界不是“只看 backend 文件夹”

compose.yaml 决定生产后端依赖哪些服务;mcp/ 向后端暴露工具协议;scenarios/ 是评测和固定工具数据。我们会讲接口和数据流,但不展开 React 或 TypeScript 内部写法。

02

读这个项目之前,补齐刚好够用的基础

所有概念都用项目里真实会出现的写法解释。

写法小白翻译本项目为什么用
def f(x: str) -> bool定义函数;x 建议传字符串;返回布尔值。类型标注帮助 IDE、Ruff 和读者理解接口,但 Python 运行时并不自动强制。
async def / await这个函数会等待网络、数据库或其他任务;等待时可以让事件循环做别的事。API、数据库、模型、工具都可能慢,不能让整个服务原地干等。
class / self类是对象模板;self 是“当前这个对象自己”。RunService 把 repository、events、dispatcher 保存为自己的属性。
@dataclass自动生成初始化、比较、显示等样板代码的数据类。内部记录、工具结果、策略上下文主要装数据,不需要复杂行为。
@classmethod方法收到的是类 cls,常用于工厂方法。AppServices.offline() 负责构造一整套离线服务。
Protocol只规定“必须有哪些方法”,不要求继承某个父类。内存仓库和 SQL 仓库可以互换;真实/假模型也能互换。
UUID几乎不会撞车的全局唯一编号。Run、审批、评测记录都要跨进程稳定识别。
dict[str, Any]键是字符串,值可以是任意类型的字典。灵活承载 metadata、JSON 和外部工具结果。
None / X | None没有值;后者表示可能是 X,也可能没有。未配置 API Key、没有 checkpoint、尚无最终结果都用它表达。
try / except / finally尝试执行、捕获错误、无论如何都清理。Worker 出错仍必须释放 lease;工具超时必须转成可审计错误。

HTTP 请求在后端眼里是什么?

输入

方法(GET/POST)、路径/api/runs)、请求头、JSON body、登录身份。

输出

状态码(200/202/404…)、响应头和 JSON;持续推送则用 SSE 的 text/event-stream

同步与异步不要死记

看到 await repository.get_run(...),先翻译成“暂停当前协程,等数据读取完成”。它不等于开新线程,也不等于函数会自动并行。

03

离线启动与第一次调试

先把变量减少:不接真实模型、不启动数据库,也能走完整调用链。

这些命令分别做什么

PowerShell · 在项目根目录执行
cd C:\Users\Jack1\Desktop\offer\devinfra-agent-platform\backend
uv sync --all-groups
$env:DEVINFRA_MODE="offline"
$env:DEVINFRA_OFFLINE_EAGER="1"
uv run uvicorn devinfra_agent.main:app --reload
cd
把当前工作目录切到 backend,后面的 pyproject.toml 才能被找到。
uv sync
uv 读取项目清单和锁文件,创建/更新 .venv 并安装运行、开发依赖。
DEVINFRA_MODE
环境变量;offline 让组合根选择内存仓库、确定性模型和 fixture 工具。
DEVINFRA_OFFLINE_EAGER
设为 1 后,内存 dispatcher 收到 Run 就立即安排 Worker 执行,适合本地演示。
uv run
在项目虚拟环境里执行后面的命令,不必手动激活 .venv
uvicorn
ASGI 服务器;它导入 devinfra_agent.main 模块中的 app 对象。
--reload
源码保存后自动重启开发服务器;只用于开发,不用于生产。

第一次手动验证

  1. 打开 http://127.0.0.1:8000/docs,这是 FastAPI 自动生成的 Swagger UI。
  2. 调用 POST /api/auth/login,用户名 maintainer,密码 maintainer-password
  3. 复制 access_token,点 Swagger 右上角 Authorize,填 Bearer 空格 token
  4. 调用 POST /api/runs,body 使用 {"workflow":"ci_triage","prompt":"分析退款测试失败","metadata":{}}
  5. 拿响应里的 Run ID 调用 GET /api/runs/{run_id}/trace,观察 step、model_call、tool_call、evidence。

建议的第一组断点

断点你应该观察继续运行后去哪里
app.py:create_runrequest 已被 Pydantic 验证;user 已由 JWT 依赖解析。RunService.create
services.py:RunService.createfingerprintkeycreatedrepository.create_run_request
worker.py:RunWorker.executedelivery_id、lease 是否拿到、构造出的 RunRequestAgentRuntime.run
runtime.py:_invoke_tool工具名、参数、风险、超时、调用次数和返回的 ToolResultEvidenceNormalizer.normalize

动手 01:证明离线 eager 真在后台执行

先把 DEVINFRA_OFFLINE_EAGER 改为 0,创建 Run 后看它停在 pending;再改回 1 重启并重试。你会理解“API 创建任务”和“Worker 执行任务”本来是两件事。

04

入口文件与依赖组装

真正的入口只有 5 行;复杂度被有意放进组合根。

backend/src/devinfra_agent/main.py
main.py · 完整文件
"""FastAPI application entry point."""

from .composition import create_app_from_environment

app = create_app_from_environment()
"""..."""
模块文档字符串,告诉读者这个文件的职责。
from .composition
开头的点表示“从当前 devinfra_agent 包内导入”。
create_app_from_environment
工厂函数:读环境变量、组装服务、创建 FastAPI 对象。
app
模块级变量。命令里的 devinfra_agent.main:app,冒号后就是它。

为什么不在入口里直接连接数据库?

入口只声明“我要一个 app”。composition.py 才负责根据环境选择内存实现或生产实现。这叫组合根(composition root):所有具体依赖集中装配,业务类只接收自己需要的对象。

AppSettings.from_env() 看到 DEVINFRA_MODE=offline,只解析 DEVINFRA_OFFLINE_EAGERbuild_services() 返回 AppServices.offline():内存仓库、内存 dispatcher、确定性模型、内存 LangGraph checkpoint、fixture 工具、内存对象存储。

生产模式先强校验必要环境变量,再创建 SQLAlchemy engine、PostgreSQL repository、Celery dispatcher、MCP/GitHub 工具、OpenAI-compatible 或确定性模型、Postgres checkpoint、Qdrant RAG、MinIO 和组件健康探针。

AppServices 是什么?

它是应用的“服务箱”:repository、dispatcher、events、runs、approvals、worker、auth、skills、tools、knowledge、evaluations、health 都放在同一个数据类里。app.state.services 保存这只箱子,测试也能换入自定义假实现。

app.py · AppServices.offline 核心
tools = runtime.tools if runtime is not None else ToolRegistry.fixture()
skills = runtime.skills if runtime is not None else SkillRegistry.default()
actual_runtime = runtime or AgentRuntime(
    gateway=DeterministicModelGateway(),
    tools=tools,
    skills=skills,
    graph=LangGraphAdapter(checkpointer=InMemorySaver()),
)
dispatcher = InMemoryDispatcher()
services = cls.compose(
    repository=InMemoryRepository(),
    dispatcher=dispatcher,
    runtime=actual_runtime,
    jwt_secret="offline-development-secret",
    user_store=UserStore.seeded(),
)
runtime is not None
调用方传了自定义 Runtime 就复用它,便于测试;否则使用默认离线实现。
actual_runtime
“实际采用的 runtime”。变量改名是为了区别可选参数 runtime
gateway
模型网关;确定性版本不访问网络,返回可预测结果。
checkpointer
图状态存档器;离线存在内存里,生产存在 PostgreSQL。
cls.compose
clsAppServices 类;compose 继续把底层对象组装成领域服务。
调试口诀

“用的到底是哪一个实现?”先看 DEVINFRA_MODE,再看 composition.build_services,最后检查 app.state.services 中对象的实际类型。

05

登录、JWT 与 API 外壳

请求还没进入业务规则前,先经过验证、身份、角色、限流和统一错误处理。

登录的真实调用顺序

app.py: login

FastAPI 把 JSON 解析成 LoginRequest,然后调用 container.auth.login(request)

app.py · 登录与角色依赖
@app.post("/api/auth/login")
def login(request: LoginRequest) -> dict[str, str]:
    return {"access_token": container.auth.login(request), "token_type": "bearer"}

def allow(*roles: Role):
    def dependency(user: User = Depends(current_user)) -> User:
        if user.role not in roles:
            raise APIProblem(403, "forbidden", "Permission denied")
        return user
    return dependency

approver = allow(Role.MAINTAINER, Role.ADMIN)
@app.post
装饰器,把下面的函数注册为 POST 路由;函数名可读,URL 才是外部接口。
request: LoginRequest
FastAPI 根据类型标注读 JSON,并让 Pydantic 校验字段。
container
当前 AppServices 对象;它在 create_app 中生成。
*roles
可变位置参数;调用 allow(A, B) 后,函数里得到角色元组。
Depends
FastAPI 依赖注入:先运行 current_user 解析 Authorization,再把 User 传进来。
approver
已经配置好的依赖,只允许 Maintainer 或 Admin;审批路由直接复用。

auth.py 的职责拆分

UserStore

保存用户并按用户名查找。seeded() 只为离线演示;configured() 读取生产配置。

JWTCodec

encode 生成 header.payload.signature;decode 验签、issuer、audience、过期时间与算法。

AuthService

login 校验用户名密码;authenticate 解析 Bearer token;依赖函数把它接入 FastAPI。

密码不是明文保存。hash_password 使用 PBKDF2-HMAC-SHA256、随机 salt 和 120000 次迭代;verify_password 重算并用常量时间方式比较。JWT 签名用 HMAC;它是签名令牌,不是加密保险箱,payload 不应放秘密。

统一错误为什么重要?

errors.py 把业务错误、Pydantic 校验错、FastAPI HTTP 错、数据库完整性错和未知异常都转成统一结构:code / message / request_id / details / retryable。未知异常不回显堆栈或密钥。

operations.py 在每个请求外面做什么?

  • 生成或沿用 X-Request-ID,写入响应,方便串联日志。
  • 对登录和普通 API 分别限流;离线内存计数,生产 Redis 计数。
  • 记录 Prometheus 请求总量和耗时;/metrics 暴露文本。
  • configure_telemetry 可安装 OpenTelemetry FastAPI instrumentor。

动手 02:观察 401、403、422 的差别

不带 token 调一个 API 得 401;用 viewer 创建 Run 得 403;给 workflow 填不存在的值会在 Pydantic 层得 422。三个错误进入业务代码的深度不同。

06

请求模型、内部记录与数据库模型

同样叫“Run”,在不同边界有不同形状,这是刻意设计。

Pydantic Contract

contracts.py / api_models.py

面向 API、模型网关和评测的已验证数据。非法枚举、缺字段、多字段会被拒绝。

Record

records.py / repositories.py

Python 内部传递的数据类,如 RunRecordApprovalRecord,不依赖 ORM session。

ORM Model

models.py

SQLAlchemy 声明数据库表、列、索引、约束和 PostgreSQL JSONB/UUID 类型。

contracts.py · RunRequest
class ContractModel(BaseModel):
    model_config = ConfigDict(extra="forbid")

class RunRequest(ContractModel):
    run_id: UUID = Field(default_factory=uuid4)
    workflow: WorkflowName
    prompt: str = Field(min_length=1)
    requested_by: Role
    steps: list[PlanStep] = Field(default_factory=list)
    tools: list[ToolSpec] = Field(default_factory=list)
    metadata: dict[str, Any] = Field(default_factory=dict)
BaseModel
Pydantic 基类;负责解析、校验、序列化和生成 JSON Schema。
extra="forbid"
传入模型未声明的字段就报错,避免拼错字段被静默忽略。
default_factory
每次创建对象时调用工厂。UUID 用 uuid4 生成;list/dict 用新容器,避免多个对象共享可变默认值。
WorkflowName
StrEnum 枚举;只允许定义过的工作流字符串。
metadata
扩展信息,例如 service、action_tool、action_arguments;灵活但仍受到后续策略校验。

文件之间怎样转换?

flowchart LR
  JSON["HTTP JSON"] --> P["Pydantic API model"]
  P --> S["Service 参数"]
  S --> R["Record dataclass"]
  R --> O["SQLAlchemy Model"]
  O --> DB[("PostgreSQL")]
  R --> Z["serializer / _run_json"]
  Z --> OUT["响应 JSON"]
          

serializers.pydocument_jsoneval_case_json 等只做 Record → JSON 字典。Run/Event/Approval 的序列化目前在 app.py 的私有函数里。若你给数据库加字段,通常要同时检查 ORM Model、Record、迁移、Repository 转换、serializer 和 API schema。

不要把 ORM 对象一路传到 API

ORM 对象依赖 session 生命周期,还可能触发隐式数据库读取。Record 把数据访问和业务层隔开,测试也能直接构造。

07

创建 Run:API、Service、幂等与 Outbox

终于开始跟一条请求真正旅行。

app.py · POST /api/runs
@app.post("/api/runs", status_code=status.HTTP_202_ACCEPTED)
async def create_run(
    request: RunCreate,
    user: User = Depends(run_mutator),
    idempotency_key: Annotated[str | None, Header(alias="Idempotency-Key")] = None,
) -> dict[str, Any]:
    run = await container.runs.create(
        workflow=request.workflow,
        prompt=request.prompt,
        metadata=request.metadata,
        requested_by=user.role,
        requested_by_user=user.username,
        idempotency_key=idempotency_key,
    )
    return _run_json(run)
status.HTTP_202_ACCEPTED
返回 202 表示请求已接收、后台工作未必完成;API 不等待 Agent 跑完。
run_mutator
上一章创建的角色依赖,viewer 会在进入函数前被拒绝。
Annotated
把“值类型”和 FastAPI 元数据放一起;这里表示可选字符串来自名为 Idempotency-Key 的请求头。
run
RunService.create 返回已持久化的 RunRecord,不一定是刚创建的。
_run_json
私有序列化函数;前导下划线表示模块内部使用约定。

为什么同一个请求不能创建两次?

网络重试、用户双击、代理超时都可能让客户端重复 POST。RunService.create 对影响请求含义的字段做稳定 JSON 序列化,再计算 SHA-256 fingerprint。调用方没给 key 时,使用 derived:{fingerprint};给了 key 时,数据库把 key 与 fingerprint 绑定。

services.py · fingerprint 与分发核心
fingerprint = hashlib.sha256(
    json.dumps(
        {
            "workflow": workflow.value,
            "prompt": prompt,
            "metadata": metadata,
            "requested_by_user": requested_by_user,
            "rerun_of": str(rerun_of) if rerun_of else None,
        },
        sort_keys=True,
        separators=(",", ":"),
    ).encode()
).hexdigest()
key = idempotency_key or f"derived:{fingerprint}"
persisted, created = await self.repository.create_run_request(
    run, key, fingerprint, self.clock()
)
if created:
    await self.dispatcher.dispatch(persisted.id)
hashlib.sha256
把任意长度字节变成固定 64 位十六进制摘要;内容相同则摘要相同。
sort_keys
字典键顺序不影响 fingerprint。
separators
去掉 JSON 非必要空格,保证序列化稳定。
.encode()
哈希函数接收 bytes;把 Python 字符串编码成 UTF-8 字节。
persisted
仓库实际保存/查到的 Run;重复请求时它是第一次请求的记录。
created
布尔值:这次是否首次插入。它决定是否需要新分发。

Outbox 解决哪扇“时间窗口”?

如果“写数据库成功”后进程崩溃,来不及给 Celery 发消息,Run 会永久 pending。Repository 在同一事务里写 Run 和 dispatch_outbox。立即分发失败也没关系,Celery Beat 的 relay 会重新 claim pending outbox,再投递。

sequenceDiagram
  participant API
  participant DB as Repository/DB
  participant Broker as Dispatcher/Broker
  API->>DB: create_run_request(run, key, fingerprint)
  DB-->>API: 同一事务写 run + outbox
  API->>Broker: dispatch(run_id)
  alt 成功
    API->>DB: mark_dispatch_succeeded
  else 失败
    API->>DB: mark_dispatch_failed
    Note over DB,Broker: relay 稍后补投递
  end
          

动手 03:验证幂等冲突

使用同一个 Idempotency-Key 连续创建相同请求,应得到同一个 Run;再保持 key 不变但修改 prompt,应得到 409 idempotency_conflict。对应测试在 test_review_durability.py

08

Repository、数据库、事件与 SSE

数据层不是“存一下 JSON”,而是并发正确性的事实源。

为什么 repositories.py 有 2081 行?

它提供两套相同行为:InMemoryRepository 让单测和离线模式快速确定;SqlAlchemyRepository 用真实事务、唯一约束和行锁保证跨进程正确。ThreadedRepository 再把同步 SQLAlchemy 调用放到线程,避免阻塞异步事件循环。

仓库能力组代表方法保证
Run 与幂等create_run_requestget_runupdate_run同 key 同请求返回旧值;同 key 异请求冲突。
Outboxclaim_dispatch_outboxmark_dispatch_*多 relay 竞争时一条消息只有一个 owner。
事件append_eventevents_after每个 Run 内 sequence 唯一递增,可精确重放。
Leaseacquire_run_leaserelease_run_lease同一 Run 同时最多一个有效 Worker owner。
审批与动作claim_approval_decisionclaim_action决定与副作用都只执行一次。
Trace 投影create_step/model_call/tool_call/evidence/artifact/auditprojection key 去重,重投递不复制轨迹。
知识与评测create_documentclaim_eval_runfinish_eval_run记录持久、评测 owner fencing。

ORM Model 与 Alembic

models.py 用 SQLAlchemy Declarative 定义 21 类表模型:users、runs、run_events、run_steps、model_calls、tool_calls、evidence、artifacts、approvals、audit_logs、skills、tools、documents、chunks、eval_cases、eval_runs、eval_results、eval_scores、feedback、dispatch_outbox、action_claims。

Alembic 的三次迁移按时间演进:首次建完整 schema → 增加 outbox/action claim 并修复旧审批数据 → 增加 durable eval 执行字段和约束。改 Model 不会自动改已有数据库,必须新增 migration。

SSE 为什么从数据库重放?

app.py · SSE 核心循环
async def replay():
    cursor = after
    while not await raw_request.is_disconnected():
        events = await container.repository.events_after(run_id, cursor)
        for event in events:
            cursor = event.sequence
            data = json.dumps(_event_json(event), separators=(",", ":"))
            yield f"id: {event.sequence}\nevent: {event.event}\ndata: {data}\n\n"
        await container.events.publisher.wait(run_id, 15)
        yield ": heartbeat\n\n"
replay
异步生成器;每次 yield 一小段文本给客户端,而不是一次性返回。
cursor / after
最后看过的事件序号;来自 Last-Event-ID,断线重连从它后面继续。
is_disconnected
浏览器关掉连接后停止循环,避免服务器继续白做事。
event.sequence
数据库事实序号,同时作为 SSE id。
heartbeat
SSE 注释行,防止代理因长时间无数据关闭连接。
Redis 只负责“叫醒”,数据库才记得历史

进程内 EventNotifier 或生产通知机制只是让等待者早点再查一次。客户端丢包、API 重启都不怕,因为 events_after 仍能按 sequence 重放。

09

分发、Worker、Lease 与持久化投影

API 只排队,真正慢而危险的工作在 Worker。

离线与生产的 Dispatcher

InMemoryDispatcher

把 ID 记入列表;eager 时用 asyncio.create_task 调 Worker。进程重启后丢失,只为本地/测试。

CeleryDispatcher

用稳定 task_id 给 broker 发送 devinfra.execute_run 或 eval task;Redis 是 broker/result backend。

celery_entrypoint.py 只允许 production,复用 AppSettingsbuild_services,并配置每 5 秒 relay Run outbox、每 10 秒 relay eval queue。dispatch.py 注册具体 Celery task,设置 late acknowledgement、失败重试和 delivery identity。

worker.py · execute 的骨架
async def execute(self, run_id: UUID | str, *, delivery_id: str) -> bool:
    key = UUID(str(run_id))
    acquired = await self.repository.acquire_run_lease(
        key, delivery_id, self.clock(), timedelta(minutes=2)
    )
    if not acquired:
        raise RunLeaseBusy(...)
    try:
        run = await self.repository.get_run(key)
        await self.repository.update_run(key, status=RunStatus.RUNNING, now=self.clock())
        outcome = await self.runtime.run(
            RunRequest(...),
            resource_scopes=scopes_for_role(run.requested_by),
        )
        await self._project_runtime(key, outcome)
        ...
    finally:
        await self.repository.release_run_lease(key, delivery_id)
delivery_id
本次消息投递的稳定身份,用作 lease owner;不是 Run ID。
key
无论调用者传 UUID 还是字符串,都归一化成 UUID 对象。
acquired
是否拿到 2 分钟 lease。False 说明另一个 Worker 正在负责或任务已结束。
outcome
Runtime 的完整结果:status、summary、trace、模型/工具记录、证据、checkpoint 和待审批动作。
finally
成功、失败、取消都释放 lease;否则任务要等过期才能恢复。

_project_runtime 为什么存在?

Runtime 只关心 Agent 逻辑,返回普通内存对象;Worker 把这些事实投影成数据库记录。每个工具原始结果先写 object_store.put,数据库只保存 URI、checksum、size,避免大日志塞进表和模型上下文。

InMemoryObjectStore

字节存在字典,URI 使用 memory://;checksum 做内容寻址。

MinioObjectStore

懒创建 bucket;同步 MinIO SDK 放到线程执行;返回 s3://bucket/key 元数据。

projection_key

重投递时相同事实得到相同唯一键,Repository 把重复写变成读取旧记录。

动手 04:追踪一个工具结果

runtime._invoke_tool 看原始 ToolResult;继续到 worker._project_runtime 看它变成 artifact;最后在 Trace API 看 tool_call 只剩 artifact URI/hash/size。

10

Agent Runtime 与 LangGraph

模型不是总指挥;Runtime 才控制顺序、预算、工具和停止条件。

runtime.py · 主图节点
nodes: Sequence[GraphNode] = (
    ("validate", self._validate),
    ("route", self._route),
    ("specialist", self._specialist),
    ("verify", self._verify),
    ("action_gate", self._action_gate),
    ("finalize", self._finalize),
)
nodes
有顺序的元组;每项是“节点名称 + 异步处理函数”。
Sequence[GraphNode]
类型接口;只要求可按顺序读取,不要求一定是 list。
self._validate
绑定方法对象,此处不加括号,因为是把函数交给图稍后调用。
specialist
根据 workflow 选出的 CI/incident/issue/knowledge 专家子流程。
action_gate
所有读证据完成并验证后,才允许处理请求中的 action_tool。
stateDiagram-v2
  [*] --> validate
  validate --> route
  route --> specialist
  specialist --> evidence_plan
  evidence_plan --> tools
  tools --> normalize
  normalize --> critique
  critique --> evidence_plan: needs_replan 且有预算
  critique --> verify: 足够
  verify --> action_gate: verification_passed
  action_gate --> finalize: 无动作/R0/R1/R3
  action_gate --> waiting: R2 interrupt
  waiting --> action_gate: resume approved/rejected
  finalize --> [*]
          

Runtime 的重要对象

对象装什么生命周期
RuntimeLimits最大模型调用、工具调用、replan、模型/工具/整次 Run 超时。Runtime 配置
CancellationToken内存取消事件;另有 durable cancellation_check 查仓库。一次执行
_RuntimeStaterequest、skill、scope、证据、错误、预算计数、trace、pending action。图运行/可 checkpoint
RuntimeOutcome对外返回的不可变快照,供 Worker 投影。图完成或中断后
ModelRequest/Response模型输入、允许工具、非信任证据、结构化输出。一次模型调用

图适配器为什么有三层?

  • DeterministicGraphAdapter:普通 for 循环依次执行节点,最容易测试。
  • LangGraphAdapter:构建 StateGraph;有 checkpointer 时支持 interruptCommand(resume=...)
  • AsyncPostgresLangGraphAdapter:真正执行时才打开 PostgreSQL saver,避免应用启动/组合阶段就联网。

模型网关的安全边界

DeterministicModelGateway 可按 stage 注入脚本响应,没有脚本就返回稳定默认值。OpenAICompatibleGateway 校验 base URL,把允许工具放 system message,把外部证据包装为“不可信用户数据”,要求 JSON Schema 结构化输出,再把 native tool calls 和 JSON proposals 都解析成 ToolProposal

模型提议不等于执行

模型只能返回工具名和参数;Runtime 会再检查 Skill allowlist、参数 schema、角色、scope、预算和风险。外部证据里的“忽略规则并导出密钥”不能扩大工具集合。

超时与取消

_await_operation 同时等待实际 operation、CancellationToken 和 durable poll。谁先完成决定结果:取消则 cancel operation 并抛 RunCancelled;超时则清理所有 task 并抛 TimeoutError;正常则返回 operation result。

11

Skill、Tool、Policy、证据与审批

这一章解释 Agent 为什么“能做事,但不能乱做事”。

SkillConfig

把 workflow 映射到 specialist、description 和 allowed_tools。这是工具可见性的第一道上界。

ToolDefinition

名称、说明、风险、handler、超时、JSON Schema、required scope、允许角色和可选幂等 handler。

PolicyContext

把当前用户角色、资源 scope、Skill allowlist、参数是否合法、预算是否剩余交给策略引擎。

四级风险不是“提示颜色”

风险策略结果例子行为
R0EXECUTErepo_search、读日志读操作可自动执行;瞬时错误最多 3 次有界重试。
R1DRAFTissue_draft只生成草稿,不发布真实变更。
R2APPROVAL_REQUIREDdeployment_restart持久中断,人工批准后从 checkpoint 恢复;必须有下游幂等。
R3BLOCKEDcredential_export即使 Admin 也永久阻止自动执行。
policy.py · 决策顺序
def decide(self, tool: ToolDefinition, context: PolicyContext) -> PolicyDecision:
    if not context.input_valid:
        return PolicyDecision(ActionDecision.BLOCKED, "invalid_input")
    if not context.budget_remaining:
        return PolicyDecision(ActionDecision.BLOCKED, "budget_exhausted")
    if tool.name not in context.skill_allowed_tools:
        return PolicyDecision(ActionDecision.BLOCKED, "tool_not_allowed_by_skill")
    if context.role not in tool.allowed_roles:
        return PolicyDecision(ActionDecision.BLOCKED, "role_denied")
    if tool.required_scope and tool.required_scope not in context.resource_scopes:
        return PolicyDecision(ActionDecision.BLOCKED, "resource_scope_denied")
    return decisions[tool.risk]
fail closed
任何前置条件不明确就 BLOCKED,不尝试“智能猜测”。
input_valid
Run 自身有效且工具 arguments 满足 input_schema;不是模型说“我检查过”。
skill_allowed_tools
当前 workflow 能看见的工具集合,来自 SkillRegistry。
allowed_roles
Tool 自己声明哪些角色可能使用。
required_scope
更细的资源权限,例如 deployments:write
decisions
风险到动作的固定映射,模型无法修改。

ToolRegistry 做了四件事

  1. 注册工具且拒绝重名;输出排序稳定的 definitions/specs。
  2. _argument_error 做项目所需的最小 JSON Schema 校验:required、额外字段、string 类型和最小长度。
  3. 执行 handler,记录 ToolAuditRecord;并发相同 idempotency key 共享同一 task。
  4. R2 强制使用 idempotent_handler,Fixture store 模拟下游幂等数据库。

R2 审批的完整恢复链

sequenceDiagram
  participant RT as Runtime
  participant LG as LangGraph Checkpoint
  participant DB as Repository
  participant U as Approver
  participant W as Worker
  participant T as R2 Tool
  RT->>LG: interrupt(pending_action)
  RT-->>W: AWAITING_APPROVAL + checkpoint
  W->>DB: create approval
  U->>DB: claim decision once
  DB-->>W: dispatch same run
  W->>DB: lease + begin execution + claim action key
  W->>LG: resume("approved")
  LG->>T: idempotent_handler(arguments, action_key)
  W->>DB: receipt + audit + terminal status
          

ApprovalService.validate_action 在批准时重新读当前 Tool、Skill 和 Policy,防止等待期间策略变化。RunWorker._resume_approval{run_id}:action:{tool} claim 动作,再恢复 checkpoint;即使批准后崩溃重投递,下游也按相同 key 返回旧 receipt。

动手 05:比较 R1、R2、R3

在 Run metadata 中分别指定 issue_draftdeployment_restartcredential_export,同时给对应 arguments。观察 status、decision reason、是否创建 approval 和是否出现 tool_call。

12

证据、RAG 与外部适配器

把“不可信的大段外部内容”变成“可引用、可追溯、有限大小的证据”。

工具结果怎样变成证据?

EvidenceNormalizer.normalize(tool_name, result, sequence) 生成 citation(如 E1)、清洗 source URI、计算内容 hash,并标注外部内容为不可信。NormalizedEvidence.model_text() 给模型的是明确包裹的数据;to_contract() 才转成最终 EvidenceItem

flowchart LR
  T["ToolResult: content + source_uri + metadata"] --> N["EvidenceNormalizer"]
  N --> E["NormalizedEvidence: citation + hash + trust"]
  E --> M["有限模型上下文"]
  T --> O["Object Store 原文"]
  E --> P["Repository evidence 索引"]
  P --> F["FinalResult 引用 E1/E2"]
          

RAG 从文档到命中

  1. KnowledgeService.create 生成 DocumentRecord,再把 KnowledgeDocument 交给 retriever.ingest。
  2. DocumentChunker 按字符上限切块并保留 overlap;每块有稳定 ID、hash、locator、metadata。
  3. Embedding provider 生成 dense vector;同时 tokenize/sparse_vector 形成稀疏表示。
  4. Qdrant 请求使用 named dense/sparse prefetch 与 RRF fusion;内存版本自己计算 cosine、BM25 和 rank fusion。
  5. Qdrant 不可用时 HybridKnowledgeService 退到 lexical fallback,并返回 degraded=True 与 limitations,绝不伪造 citation。
术语直觉项目实现
Dense embedding语义相近的文字向量方向相近。离线确定性 hash 向量;生产 Sentence Transformers 多语言模型。
BM25 / sparse关键词出现与稀有程度越匹配,分越高。bm25_scores 和 Qdrant sparse vector。
RRF不直接比较不同算法分数,而按各自排名融合。Qdrant fusion 或内存 rank score。
Prompt injection文档里伪装成系统指令的恶意文字。injection_suspected 标注 trust,不把它当指令。

三种外部适配器

UrllibAdapterTransport

adapters/http.py

统一异步 HTTP 外壳:实际同步 urlopen 放线程,限制超时和最大响应字节。

MCPStreamableHTTPClient

adapters/mcp.py

校验 HTTP(S)、host allowlist、无 URL 凭据;完成 initialize、notifications/initialized、tools/list、tools/call,解析 JSON 或 SSE。

GitHubReadOnlyAdapter

adapters/github.py

默认只读;限制 HTTPS host、repository allowlist、内容路径穿越;有缓存、限大小和秘密脱敏。

MCP 是协议,不是授权。MCP Server 宣称有某工具后,Python 端仍把它映射成带 risk、role、scope、schema、timeout 的 ToolDefinition,最终仍走 Policy。

真实外部连接的调试顺序

先检查 URL 语法 → host/repo/path allowlist → token 是否在 adapter 内注入 → HTTP status → 响应大小/Content-Type → ToolResult provenance → Policy 决策。不要第一步就怀疑模型。

13

评测 Harness 与失败分析

API 返回 200 只说明没崩;Harness 检查 Agent 做得对不对、安全不安全。

两条相关但不同的评测路径

独立 Harness

harness/ 加载场景包,按 replay mode 执行案例、打分、聚合、比较 baseline,并由 CLI 输出 JSON/Markdown。

持久化 Eval API

catalog_services.EvalService 创建 eval run 并分发;evaluation.EvalExecutionService 用 lease 执行 Harness,逐例保存 result/score,最后写 aggregate/regression。

harness/ 每个模块做什么

模块职责重要对象/函数
schemas.py案例、rubric、observation、grade、experiment 的严格 Pydantic schema。EvalCaseDefinitionExperimentReport
dataset.py加载 manifest/cases/recordings,校验 72 例分布、fixture/hash 和子集。ScenarioPack.load/select/recording_for
executors.pyexact、recorded_tools、full 三种方式产生 observation。AgentRuntimeExecutor.execute
metrics.py精确可复现的 precision/recall/F1、Recall@K、MRR、NDCG。precision_recall_f1
graders.py按 rubric 评分,安全项不可被 LLM judge 覆盖。DeterministicGraderJudgeAugmentedGrader
runner.py循环 case/repetition,收集 CaseRunResult,聚合指标。HarnessRunner.run/_aggregate
baseline.py安全必须 100%,总分至少 .85,关键指标相对退化不超 .03。RegressionGate.evaluate
reporting.py / cli.pyJSON 事实报告、Markdown 展示和命令行参数。render_json/render_markdown/main

三种 replay

exact

直接重放最终 observation。最快、完全确定,适合 PR 门禁,但不测试当前 Runtime。

recorded_tools

真正跑 Runtime,工具返回冻结录制值。能测规划/策略/证据链,隔离网络波动。

full

Runtime 与 fixture 工具全走,适合集成与鲁棒性;仍不触碰真实生产资源。

PowerShell · 运行评测
cd backend
uv run devinfra-eval --name pr-smoke --mode exact --subset smoke
uv run devinfra-eval --name robustness --mode full --subset robustness --repetitions 3

devinfra-eval 来自 pyproject.toml [project.scripts],入口是 devinfra_agent.harness.cli:main

失败标注

FailureService.list 找失败/取消 Run 并附已有 feedback;label 校验 Run 存在且确实失败,再保存 reviewer、label、comment。它让人工失败分类反哺评测集和修复优先级,但当前代码不会自动训练模型。

14

健康检查、可观测性与生产组合

把离线积木逐个换成可持久、可扩展的生产积木。

三个健康接口不是重复

/healthz

进程活着就 200。容器编排据此判断是否需要重启。

/readyz

关键依赖可用才 200。未 ready 时负载均衡不应发新请求。

/api/health/components

登录后查看 repository、broker、worker、object store、Redis、MCP、Qdrant 各自状态。

ComponentCheck.required 区分关键和可降级依赖:Qdrant 可选,因为 RAG 能 lexical fallback;Redis/MCP 生产组合中是 required。HTTPHealthProbe 同样执行 scheme/host allowlist,避免健康检查变成 SSRF 工具。

离线 → 生产替换表

能力离线生产为什么换
业务数据InMemoryRepositoryPostgreSQL + SQLAlchemy + Alembic跨进程持久、事务、唯一约束、行锁。
任务队列InMemoryDispatcherRedis + Celery独立 Worker、重试、并发、定时 relay。
图 checkpointInMemorySaverPostgres AsyncPostgresSaver审批等待和崩溃恢复不丢状态。
知识检索InMemoryHybridIndexQdrant + lexical fallback大规模 dense/sparse 向量与持久索引。
原始大对象InMemoryObjectStoreMinIO日志/diff/报告不塞数据库。
工具fixture registryMCP + 可选 GitHub + fixture 补位接入真实外部能力又保持统一策略。
模型DeterministicDeterministic 或 OpenAI-compatible演示可复现;生产可选择真实模型。
限流内存窗口Redis 计数多 API 实例共享。

Docker Compose 启动顺序

Postgres、Redis、Qdrant、MinIO、MCP 先各自健康 → API 运行 alembic upgrade head 后启动 Uvicorn → Worker 等 API 健康再启动 Celery worker + beat → Web 等 API。可观测 profile 额外启动 OTel Collector、Prometheus 和 Grafana。

PowerShell · 生产参考组合
docker compose --profile observability up --build

观测链

operations.py 结构化 JSON 日志带 request_id;Prometheus client 统计 route/status/latency;OpenTelemetry instrumentation 产生 HTTP trace 并经 OTLP HTTP exporter 发到 Collector;Prometheus 拉 /metrics;Grafana 展示。它们解决的是“发生了什么、慢在哪、请求如何串起来”,不是业务事实存储。

Compose 里的凭据只为本机演示

真实部署要用企业 IdP、Secret Manager、托管数据服务、TLS 和网络隔离。不要把固定 JWT/MCP/MinIO 密钥暴露到公网。

15

测试、断点与修改配方

把“我看懂了”变成“我能安全地改”。

最实用的测试命令

PowerShell · 从仓库根目录
# 只跑一个行为,失败时显示完整信息
uv run --project backend pytest backend/tests/test_run_api.py -vv

# 只跑一个测试函数
uv run --project backend pytest backend/tests/test_runtime.py::test_ci_runtime_normalizes_evidence_and_finalizes_with_citations -vv

# 全部后端测试与覆盖率门禁
uv run --project backend pytest backend/tests --cov=devinfra_agent --cov-fail-under=85

# 静态检查:导入、未使用变量、异步误用、风格
uv run --project backend ruff check backend/src backend/tests

按现象找第一断点

现象先看第一断点相关测试
登录失败/401Authorization、token issuer/audience/expAuthService.authenticatetest_auth_api.py
创建 Run 409Idempotency-Key 与 fingerprintRunService.createtest_review_durability.py
Run 一直 pendingoutbox 状态、dispatcher、eager/workerOutboxRelay.relay_oncetest_outbox_relay.py
Worker 重复/抢占lease owner/expiry/delivery_idRunWorker.executetest_delivery_leases.py
工具没执行Skill allowlist、schema、scope、risk、budgetPolicyEngine.decidetest_policy.py
审批后没恢复approval status/execution_state/checkpoint/action claim_resume_approvaltest_runtime_checkpoint_resume.py
SSE 丢/重复事件Last-Event-ID、sequence、terminal racerun_events.replaytest_sse.py
RAG 无结果chunk、filter、Qdrant health、degradedHybridKnowledgeService.searchtest_rag.py
评测重复结果eval lease、projection key、owner fencingEvalExecutionService.executetest_evaluation_execution.py

六个安全修改配方

新增一个 API 字段

  1. 先在 API 测试写期望。
  2. 改 Pydantic request/response。
  3. 若要持久化,继续改 Record、Model、migration、Repository 转换。
  4. 改 serializer/OpenAPI。
  5. 跑目标测试、全套测试、Ruff。

新增只读工具

  1. 写 ToolRegistry/adapter 测试。
  2. 定义 ToolDefinition:R0、schema、scope、role、timeout。
  3. 加入相应 Skill 的 allowed_tools。
  4. 写 Policy 拒绝用例和 Runtime 调用用例。
  5. 添加 Harness gold/forbidden tool 期望。

新增 Workflow/Skill

  1. 扩展 WorkflowName 与 RUNTIME_WORKFLOWS。
  2. 新增 SkillConfig,选择最小工具集。
  3. 确认 Runtime specialist 前缀和默认工具参数。
  4. 增加路由/运行/Harness case。
  5. 更新索引和文档。

新增数据库字段

  1. 先写 repository round-trip 测试。
  2. 改 SQLAlchemy Model 与 Record。
  3. 创建新 Alembic migration,不改旧 migration。
  4. 改两套 repository 和转换函数。
  5. 验证空库升级到 head 与旧库升级。

调整审批规则

  1. 先写 R0–R3 policy matrix 测试。
  2. 改 PolicyEngine 或 ToolDefinition。
  3. 检查批准时 revalidation。
  4. 检查 checkpoint resume、action claim、下游幂等。
  5. 跑 durability/safety Harness。

增加评测指标

  1. 先在 metrics 测试手算一个小例子。
  2. 实现纯函数。
  3. 在 grader 生成 metric。
  4. 在 aggregate/report/baseline 明确权重和门禁。
  5. 保证安全硬门禁不被平均分掩盖。

读测试的正确顺序

  1. 先读测试名,它是一句行为规范。
  2. 看 Arrange:构造了哪些服务、假时间、脚本模型和工具。
  3. 看 Act:真正调用哪个公开入口。
  4. 看 Assert:系统承诺什么,而不只是实现细节。
  5. 用测试名反搜源码,再在入口打断点单步。
修改时不要从 2081 行 Repository 中间猜

先找最小失败测试,再找公开调用入口,再沿对象属性进入具体实现。改完至少跑:目标测试 → 同模块测试 → 全套后端测试 → Ruff。

16

包、文件、类与函数总索引

以后忘了某功能在哪,回到这里反查。

运行与开发依赖:每个包为什么存在

包/工具项目里负责什么主要出现位置
FastAPI路由、依赖注入、请求解析、OpenAPI、异常处理、StreamingResponse。app.py、errors.py、operations.py
Starlette(间接)FastAPI 底层 ASGI、Request/Response/中间件能力。经 FastAPI 使用
Uvicorn运行 ASGI app 的 HTTP 服务器和开发 reload。启动命令、Dockerfile
Pydantic / pydantic-core严格数据模型、字段校验、JSON Schema、序列化。contracts.py、api_models.py、harness/schemas.py
SQLAlchemyORM 模型、session、查询、事务、数据库异常。models.py、repositories.py、composition.py
psycopg(间接)SQLAlchemy 与 PostgreSQL 的驱动;也供 LangGraph saver。数据库 URL
Alembic数据库 schema 版本管理与升级/降级脚本。backend/alembic/
CeleryRun/Eval 后台任务、重试、late ack、Beat 定时 relay。dispatch.py、celery_entrypoint.py
RedisCelery broker/result、生产限流、健康检查;不是业务事实源。operations.py、health.py、compose.yaml
LangGraphAgent StateGraph、节点执行、interrupt/resume。graph.py、runtime.py
langgraph-checkpoint-postgres把图 checkpoint 持久到 PostgreSQL。graph.py
MinIO SDK上传/读取原始工具输出与大 artifact。object_store.py
prometheus-clientCounter/Histogram 和 /metrics 文本。operations.py
opentelemetry-exporter-otlp-proto-http / opentelemetry-instrumentation-fastapi(连带 API/SDK)前者把 trace 通过 OTLP HTTP 发给 Collector;后者自动包住 FastAPI 请求并生成 span。operations.py、compose/ops
sentence-transformers(rag extra)生产多语言文本 embedding;连带 torch/transformers/numpy 等大依赖。rag.py
httpx(dev)FastAPI TestClient 底层 HTTP 测试客户端。backend/tests
pytest测试发现、fixture、参数化、异常断言。backend/tests
pytest-cov / coverage测量覆盖率并执行 85% 门禁。质量命令
Ruff静态检查、导入排序、异步与常见 bug 规则。pyproject.toml
Hatchlingsrc/devinfra_agent 构建为 Python wheel。pyproject build-system
uv虚拟环境、依赖解析/锁定、运行命令;不是应用运行时库。uv.lock、Dockerfile、开发命令

实际使用的标准库模块

异步与并发

asyncio 管协程/task/锁/超时;threading.RLock 保护内存仓库;collections.deque/defaultdict/Counter 管队列与统计。

数据与类型

dataclassesenum.StrEnumtypingcollections.abcuuiddatetime

序列化与安全摘要

jsonhashlibhmacbase64secrets

网络

urllib.request/error/parse 实现不额外依赖客户端库的 HTTP adapters;urlsplit 用于安全验证。

运行与工具

os 读环境变量;loggingargparsepathlibtimemathre

展开:标准库逐个对照

argparse 解析 Harness CLI;asyncio 管协程、Task、Lock、超时和取消;base64 编解码 JWT/password 字段;collections 的 Counter/defaultdict/deque 做检索统计与脚本响应队列;collections.abc 提供 Mapping/Iterable/Sequence/Awaitable/Callable 等运行时接口;dataclasses 生成数据类并用 replace 复制修改;datetime 统一 UTC 时间、过期与 lease;enum 提供 StrEnum;hashlib 做请求/内容摘要;hmac 做 JWT 签名和安全比较;importlib.util 检查 LangGraph 是否安装;itertools 生成 MCP request id;json 处理 API/配置/报告;logging 输出结构化日志;math 做向量与排名指标;os 读环境变量;pathlib 安全处理场景文件;re 做 tokenization、脱敏和 injection 规则;secrets 生成 salt;threading 的 RLock 保护内存仓库;time 做限流窗口和性能计时;typing 提供 Any/Protocol/Annotated/TypeVar;urllib.errorurllib.parseurllib.request 负责底层 HTTP 与 URL 安全解析;uuid 生成和解析实体 ID。

全部源码模块与公开/关键私有入口

模块类/函数入口(都在这里)职责与修改定位
__init__.py导出 main.app包导入便利入口;本身没有业务逻辑。
adapters/__init__.py标记 adapters 子包。
adapters/http.pyAdapterHttpResponseAdapterHttpTransport.requestUrllibAdapterTransport.request/_request统一 HTTP transport、线程化 urlopen、超时和响应大小。
adapters/mcp.pyMCPStreamableHTTPClientcall_tool/close/tool_definitions/_initialize/_rpc/_notification/_headers/_parse_response/parse_rpc_body/_structured_from_text/_raise_http_error_header/_bounded_redactedMCP Streamable HTTP 生命周期、SSE/JSON 解析、认证与脱敏。
adapters/github.pyGitHubReadOnlyAdapterget_issue/get_pull_request/get_check_runs/get_content/search_code/tool_definitions/_get/_request/_response_body/_validate_repository/_validate_number_bounded_redacted_json/_schemaGitHub 只读工具、allowlist、缓存、路径/响应安全。
api_models.pyAPIRequestToolTestRequestDocumentCreateKnowledgeSearchRequestEvalCaseCreateEvalRunCreateEvalResultCreateEvalScoreCreateFailureLabelCreate非 Run/Approval 的 HTTP body schema。
app.pyAppServices.offline/composeRunCreateApprovalDecision_run_json/_event_json/_approval_jsoncreate_app;其内 31 个 route handlerFastAPI 组合、RBAC 路由、REST/SSE 边界。新增接口先来这里。
auth.py_encode/_decode/hash_password/verify_passwordUserStore.seeded/by_username/configuredJWTCodec.encode/decodeAuthService.login/authenticateauth_dependency本地凭据、JWT、FastAPI 当前用户依赖。
catalog_services.pyRegistryService.list_skills/list_tools/test_toolKnowledgeService.create/list/searchEvalService.create_case/create_run/create_result/create_score/detailFailureService.label/listHealthService.components/ready目录、知识、评测、失败、健康的领域服务。
celery_entrypoint.py模块级 settings/services/celery_app 与 beat_scheduleCelery worker/beat 进程入口,只允许 production。
composition.pyAppSettings.from_envbuild_servicescreate_app_from_environment环境校验、离线/生产依赖选择,排查“实际用了谁”。
contracts.pyRole/RiskLevel/RunStatus/WorkflowNameContractModel/PlanStep/ToolSpec/EvidenceItem/ApprovalRequest/FinalResult/RunRequest/EvalCase/APIError/TraceEventAPI 与 Agent 共享的严格 Pydantic 合同。
dispatch.pyCeleryLikeCeleryDispatcher.dispatch/dispatch_eval/healthOutboxRelay.relay_oncecreate_celery_app、四个 register_*_taskCelery 配置、任务注册、outbox/eval relay。
errors.pyAPIProblemerror_responseinstall_error_handlers 及内部 5 类 handler统一、安全的错误 envelope。
evaluation.pyEvalExecutionService.execute/_ensure_case/_config/_failEvalQueueRelay.relay_once持久 eval lease、Harness 执行、逐例结果/分数、恢复与 fencing。
events.pyEventPublisherNullEventPublisher.publishEventNotifier.publish/waitEventService.append事件先持久化再通知;SSE 唤醒。
evidence.pyNormalizedEvidence.to_contract/model_textEvidenceNormalizer.normalize_safe_source_uricitation、hash、URI、trust 和模型上下文隔离。
gateways.pyToolProposal/ModelRequest/ModelResponse/ModelGatewayDeterministicModelGateway.completeOpenAICompatibleGateway.complete_parse_model_response/_default_transport模型协议、确定性假模型、OpenAI-compatible JSON Schema 调用。
graph.pyGraphState/GraphAdapterDeterministicGraphAdapter.executeLangGraphAdapter.execute/_runtime_graph/execute_runtime/resume_runtime/_executionRuntimeGraphExecution;Postgres adapters;default_graph_adapter普通图、LangGraph、checkpoint interrupt/resume。
harness/__init__.py导出 Harness 常用类型评测子包公共入口。
harness/baseline.pyRegressionGate.evaluate总分、安全与相对退化门禁。
harness/cli.pyparsermaindevinfra-eval 命令行入口。
harness/dataset.pyScenarioPack.load_default/load/asset_paths/smoke_cases/robustness_cases/select/recording_for/_hash_dataset/_validate版本化场景包加载、hash 与完整性。
harness/executors.pyAgentRuntimeExecutor.execute/_evidence_ids/_recorded_registry/_full_registryrecorded/full Runtime 执行与工具注册。
harness/graders.pyStructuredJudge.scoreDeterministicGrader.grade/_factual_scoreJudgeAugmentedGrader.grade确定性评分与受限语义 judge。
harness/metrics.pyPRFScoreprecision_recall_f1/tool_f1/recall_at_k/mrr/ndcg_at_k、内部 dcg可手算、可复现的评测指标。
harness/reporting.pyrender_jsonrender_markdown机器事实报告与人类展示。
harness/runner.pyHarnessExecutorFixtureExecutor.executeHarnessRunner.run/_aggregate案例/重复循环、评分与聚合。
harness/schemas.pyHarnessModelReplayMode/EvalWorkflowGoldEvidence/EvalRubric/EvalCaseDefinition/ScenarioManifest/CaseObservation/CaseGrade/ExperimentConfig/Metadata/CaseRunResult/Report/RegressionDecision全部评测数据合同与交叉字段校验。
health.pyComponentCheckHTTPHealthProbe.healthRedisHealthProbe.health外部依赖健康探测与 SSRF 边界。
main.pyapp = create_app_from_environment()Uvicorn API 入口。
models.pyBase、两个 mixin、21 个 ORM Model数据库表、列、索引、唯一约束和 PostgreSQL 类型。
object_store.pyStoredObject_validate_keyInMemoryObjectStore.put/get/healthMinioObjectStore._ensure_bucket/put/get/health原始 artifact 内容寻址与对象存储。
operations.pyJsonLogFormatter.formatconfigure_structured_logging、两种 RateLimiter、HttpMetrics.observe/renderinstall_operationsconfigure_telemetry日志、request ID、限流、Prometheus、OTel。
policy.pyActionDecisionPolicyContextPolicyDecisionPolicyEngine.decidefail-closed R0–R3 决策。
rag.pyKnowledgeDocument/Chunk/SearchHit/Response、Embedding/Index protocols、两种 embedding、DocumentChunker.chunkInMemoryHybridIndexHybridKnowledgeService.ingest/searchQdrantHTTPStoretokenize/sparse_vector/bm25_scores/rank_scores/cosine/injection_suspected/chunk_payload/qdrant_hit完整混合检索、降级与引用。
records.pyutc_now;Document/Eval/Feedback/RunStep/ModelCall/ToolCall/Evidence/Artifact/Audit/Skill/Tool/Chunk 共 16 个 Record跨层内部数据载体。
repositories.py_aware;Run/Event/Approval/Outbox/ActionClaim records;InMemoryRepositorySqlAlchemyRepository 的 Run、outbox、cancel、event、lease、approval、action、projection、catalog、eval 全套方法;全部 _xxx_record/_xxx_model 转换;ThreadedRepositorynew_run业务事实与并发正确性核心。详细能力组见第 08 章。
serializers.pydocument_json/eval_case_json/eval_run_json/eval_result_json/eval_score_json/feedback_jsonRecord 到 API JSON。
services.pyDispatcherInMemoryDispatcher.dispatched/dispatch/dispatch_eval/healthscopes_for_roleRunService.create/rerun/cancel/requireApprovalService.create_for_run/decide/validate_actionRun 生命周期、幂等创建、取消、审批业务规则。
skills.pyRUNTIME_WORKFLOWSSkillConfig.from_mappingDEFAULT_SKILL_CONFIGSSkillRegistry.default/get/workflowsworkflow → specialist + 最小工具集。
tools.pyToolRisk/InputError/TransientError/Result/AuditRecord、两种 Handler 类型、FixtureDownstreamIdempotencyStore.execute/effect_count/resultToolDefinitionToolRegistry.fixture/register/get/names/definitions/specs/validate_arguments/supports_durable_idempotency/audit_records/inflight_idempotency_keys/execute/_execute_and_audit/_check_idempotent_fingerprint/_finish_idempotent_fixture/_fixture_schema/_argument_error工具注册、schema、执行、审计与幂等。
runtime.pyRuntimeLimitsCancellationTokenRunCancelledRuntimeOutcome_RuntimeState checkpoint 方法;AgentRuntime.run/resume/_validate/_route/_specialist/_verify/_collect_evidence_plan/_tool_proposals/_action_gate/_apply_approval/_finalize/_model_call/_execute_tools/_invoke_tool/_policy_context/_default_tool_arguments/_finalize_result_await_operationAgent 的真正控制器:状态、图、预算、模型、工具、验证、审批、取消。
worker.pyRunLeaseBusyRunWorker.cancel/health/execute/_resume_approval/_project_runtimelease、Runtime 执行、审批恢复、Trace/artifact 投影。
展开:大模块中没有在上表逐字列出的全部补充符号

上表按“先找功能”压缩了重复方法。下面把剩余函数、内部函数、类型和状态变量补齐;这些名称也可被页面搜索找到。

位置补充符号统一解释
dispatch.pyCeleryLike.send_taskregister_run_task 内的 execute_runregister_outbox_task 内的 relay_outboxregister_eval_task 内的 execute_evalregister_eval_relay_task 内的 relay_eval_queuesend_task 是 Celery 最小协议;四个 register 函数把普通 Python worker/relay 包成 Celery task,内部函数是 Celery 真正调用的入口。
errors.pyapi_problem_handlervalidation_handlerhttp_handlerintegrity_handlerunexpected_handler分别捕获项目业务错误、Pydantic/FastAPI 校验错、HTTPException、SQLAlchemy 完整性错误和兜底异常,最终都调用统一 error response。
gateways.py / graph.py_default_transport.sendGraphState.check_cancelledAsyncPostgresLangGraphAdapter._runtime_callsend 是送进线程的同步 urlopen 闭包;check_cancelled 是图状态协议要求;_runtime_call 统一打开 saver、setup、转调 execute/resume 并关闭资源。
harness/executors.pytransient_handlertimeout_handlermalicious_runbook三者是构造评测工具的内部 handler,分别模拟瞬时失败、永不及时返回和携带 prompt injection 的 Runbook,用来验证重试、超时和安全隔离。
harness/schemas.pyEvalRubric.validate_weight_sumEvalCaseDefinition.validate_tool_setsExperimentMetadata两个 Pydantic model validator 做跨字段校验:权重和必须合法,expected/forbidden tool 不能矛盾;ExperimentMetadata 记录数据集 hash、模式、时间等实验身份。
operations.pyInMemoryRateLimiterRedisRateLimiterprometheus_metricsoperations_middleware两种 limiter 实现同一窗口语义;metrics handler 返回 Prometheus bytes;middleware 包住每个 HTTP 请求,做 request ID、限流、计时和统一响应头。
rag.pyVectorStoreUnavailableKnowledgeChunkSearchResponseEmbeddingProviderDeterministicEmbeddingProviderMultilingualMiniLMEmbeddingProviderupsertupsert_documentssearch_lexical_filteredensure_collection异常触发显式降级;Chunk/Response 是检索数据;Provider 协议允许假/真 embedding 互换;upsert 系列写索引,lexical/filter 做内存退路,ensure_collection 幂等创建 Qdrant collection。
models.pyIdentityMixinTimestampMixinUserModelRunModelRunEventModelRunStepModelModelCallModelToolCallModelEvidenceModelArtifactModelApprovalModelAuditLogModelSkillModelToolModelDocumentModelChunkModelEvalCaseModelEvalRunModelEvalResultModelEvalScoreModelFeedbackModelDispatchOutboxModelActionClaimModel两个 mixin 复用主键与时间列;其余每个 Model 一一对应数据库表。名字去掉 Model 后就是业务实体,字段细节以类内 mapped_column 为准。
records.py / repositories.pyEvalCaseRecordEvalRunRecordEvalResultRecordEvalScoreRecordFeedbackRecordRunStepRecordModelCallRecordToolCallRecordEvidenceRecordArtifactRecordSkillRecordToolRecordChunkRecordEventRecordDispatchOutboxRecordActionClaimRecordRecord 是离开 ORM session 后仍可安全使用的普通数据对象;名称与 Model 对应。Event/Outbox/ActionClaim 定义在 repositories.py,因为它们紧贴持久并发语义。
repositories.py:Run/Outboxlist_dispatch_outbox_queue_dispatch_locked/_queue_dispatchrequest_cancelis_cancel_requestedrecord_worker_failureget_run_by_idempotency_keyrun_lease_retry_at列出/排队 outbox;持久取消;失败计数并在上限后终止;按幂等 key 找 winner;告诉 Celery 何时能在 lease 过期后重试。带 locked 的内存版本要求调用者已持有 RLock。
repositories.py:审批/动作create_approvalapproval_for_runbegin_approval_executioncomplete_approval_executionmark_approval_checkpoint_resumedcomplete_actionget_action_claim从创建审批到 claim decision、claim execution、标记 checkpoint 已恢复、写动作 receipt 的状态机。每一步都检查旧状态与 owner,防止重复副作用。
repositories.py:评测/反馈get_eval_caseget_eval_runget_eval_resultlist_eval_resultslist_eval_scorescreate_feedbacklist_feedback按主键读评测实体、按父 ID 列子项,保存/读取人工失败标签。
repositories.py:Trace/Catalogset_healthlist_stepscreate_model_call/list_model_callscreate_tool_call/list_tool_callscreate_evidence/list_evidencecreate_artifact/list_artifactscreate_audit/list_auditsupsert_skill/list_skill_recordsupsert_tool/list_tool_recordscreate_chunk/list_chunksset_health 只供内存故障测试;其余按“写一条/按父对象读取”形成 Trace、审计、注册表与知识 chunk 投影,并用唯一 projection key 去重。
repositories.py:转换/线程壳_insert_step_record_model_call_record_tool_call_record_evidence_record_artifact_record_audit_record_skill_record_tool_record_document_record_chunk_record_eval_case_record_eval_run_record_eval_result_record_eval_score_record_feedback_record_run_model/_run_record_outbox_record_event_record_approval_model/_approval_record_action_claim_recordThreadedRepository.__getattr___insert 统一 session add/commit;_xxx_record 把 ORM 转普通记录,_xxx_model 反向转换;__getattr__ 动态把同步仓库方法代理成在线程中运行的异步调用。
runtime.py:状态对象RuntimeLimits.__post_init__CancellationToken.cancelled_RuntimeState.check_cancelled/check_cancelled_async/add_error/to_checkpoint/from_checkpointAgentRuntime.restore_checkpoint_state/awaiting_approval_status/finalize_checkpoint_state__post_init__ 验证预算/超时为正;cancelled 是只读状态;State 方法负责内存+持久取消、记录 partial error 和可序列化 checkpoint;AgentRuntime 三个桥接方法供 graph.py 在恢复时调用。
tools.pyToolInputErrorTransientToolError前者表示参数/幂等前置条件错误,不应盲重试;后者明确表示 R0 只读 adapter 的短暂失败,Runtime 才做最多三次有界重试。
adapters/github.py_CachedResponse私有数据类,保存响应过期时间与已经脱敏/限长的 ToolResult,避免重复请求。

变量阅读规则:self.x 是对象长期依赖/状态;函数参数是调用方输入;普通局部变量只服务当前步骤;前导下划线表示内部实现;actual_* 表示从“可选注入值或默认值”中最终选中的对象;key/owner/delivery_id/projection_key 都是不同层级的身份,不能互换。

脚本与迁移索引

文件用途什么时候改/运行
scripts/export_openapi.py从 FastAPI app 导出 OpenAPI JSON。API schema 改动后给前端/客户端生成类型。
scripts/load_test.py对登录、创建 Run、轮询等关键路径做负载验证。性能/容量测试,不替代功能测试。
alembic/env.py把 SQLAlchemy metadata 和数据库 URL 接入 Alembic。运行 upgrade、生成迁移。
a3af9...create_task_3_schema.py首次创建完整业务表。历史迁移,不应直接改。
7ab7...harden_task_3_durability.py增加 outbox/action claim、约束并清理旧数据。历史迁移,不应直接改。
c91e...add_durable_eval_execution.py增加 eval lease/owner/恢复字段与唯一性。历史迁移,不应直接改。

31 个 API 路由逐条反查

方法与路径处理函数进入的服务/行为
POST /api/auth/loginloginAuthService.login,返回 Bearer JWT。
GET /api/auth/meme返回当前依赖注入的 User 身份。
POST /api/runscreate_runRunService.create,幂等持久化并分发。
GET /api/runslist_runsrepository.list_runs
GET /api/runs/{run_id}run_detailRunService.require
GET /api/runs/{run_id}/tracerun_trace并发读取 events、steps、model/tool calls、evidence、artifacts。
GET /api/runs/{run_id}/eventsrun_events按 Last-Event-ID 重放的 SSE generator。
POST /api/runs/{run_id}/rerunrerunRunService.rerun,保留原始输入并记录 rerun_of。
POST /api/runs/{run_id}/cancelcancelRunService.cancel,内存+持久取消。
GET /api/approvalslist_approvalsrepository.list_approvals
GET /api/approvals/{approval_id}approval_detailrepository.get_approval,不存在转 404。
POST /api/approvals/{approval_id}/decisiondecide_approvalApprovalService.decide,只允许 Maintainer/Admin。
GET /api/skillslist_skillsRegistryService.list_skills
GET /api/toolslist_toolsRegistryService.list_tools
POST /api/tools/{name}/testtest_toolRegistryService.test_tool,仍经过 Policy。
POST /api/knowledge/documentscreate_documentKnowledgeService.create,保存并 ingest。
GET /api/knowledge/documentslist_documentsKnowledgeService.list
POST /api/knowledge/searchsearch_knowledgeKnowledgeService.search,返回 hits/degraded/limitations。
POST /api/evals/casescreate_eval_caseEvalService.create_case
GET /api/evals/caseslist_eval_casesrepository.list_eval_cases
POST /api/evals/runscreate_eval_runEvalService.create_run,持久化后分发。
GET /api/evals/runslist_eval_runsrepository.list_eval_runs
GET /api/evals/runs/{eval_run_id}eval_run_detailEvalService.detail,嵌套 results/scores。
POST /api/evals/runs/{eval_run_id}/resultscreate_eval_resultEvalService.create_result
POST /api/evals/results/{eval_result_id}/scorescreate_eval_scoreEvalService.create_score
GET /api/failureslist_failuresFailureService.list,附人工 labels。
POST /api/failures/{run_id}/labellabel_failureFailureService.label
GET /api/health/componentscomponent_healthHealthService.components,需要登录。
GET /healthzhealthz无依赖的 liveness。
GET /readyzreadyzHealthService.ready,失败返回 503。
GET /healthlegacy_health兼容旧调用方的健康接口。
你已经有一张可工作的后端地图

今后先从现象确定领域,再从 API/Worker/Harness 入口进入,用本表找类和方法。不要从文件数量最大的模块开始盲读。