第一遍:只走主链
从第 01 章读到第 11 章。先接受“一个请求会经过很多层”,不要急着记每个方法。
main.py 出发,亲手走完一个企业级 Agent 后端这不是把名词堆给你的架构说明。你会跟着一条真实请求,依次看懂入口、API、Service、Repository、Worker、Agent、工具、证据、审批、RAG、评测与生产设施,最后能自己定位、调试和修改代码。
你的目标不是背代码,而是形成“看到现象 → 找到入口 → 沿调用链定位”的能力。
从第 01 章读到第 11 章。先接受“一个请求会经过很多层”,不要急着记每个方法。
照第 03、15 章启动离线服务,在 create_run、RunService.create、RunWorker.execute 处打断点。
从修改配方挑一个任务,先找测试,再改实现。此时第 16 章是你的后端地图。
先看它旁边的“新手补课”和变量表。页面顶部可搜索任何函数、变量、包或接口。代码里的 self、await、Depends 等只在第一次出现时完整解释,后面会链接回来。
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,最后把可审计结果写回仓库”。
先知道每个房间放什么,再进房间看家具。
整个仓库。里面既有 Python 后端,也有前端、MCP 服务、场景数据和运维文件。
devinfra_agent 是可导入的 Python 包;目录里有 __init__.py。
一个 .py 文件就是模块,例如 services.py;用 from .services import RunService 引入。
src/ 布局?源码不直接躺在 backend/ 根目录,而在 backend/src/devinfra_agent。这样测试必须像真实安装后的用户一样导入包,可避免“只因当前目录碰巧在搜索路径里所以能运行”的假成功。
compose.yaml 决定生产后端依赖哪些服务;mcp/ 向后端暴露工具协议;scenarios/ 是评测和固定工具数据。我们会讲接口和数据流,但不展开 React 或 TypeScript 内部写法。
所有概念都用项目里真实会出现的写法解释。
| 写法 | 小白翻译 | 本项目为什么用 |
|---|---|---|
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;工具超时必须转成可审计错误。 |
方法(GET/POST)、路径(/api/runs)、请求头、JSON body、登录身份。
状态码(200/202/404…)、响应头和 JSON;持续推送则用 SSE 的 text/event-stream。
看到 await repository.get_run(...),先翻译成“暂停当前协程,等数据读取完成”。它不等于开新线程,也不等于函数会自动并行。
先把变量减少:不接真实模型、不启动数据库,也能走完整调用链。
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
backend,后面的 pyproject.toml 才能被找到。uv 读取项目清单和锁文件,创建/更新 .venv 并安装运行、开发依赖。offline 让组合根选择内存仓库、确定性模型和 fixture 工具。.venv。devinfra_agent.main 模块中的 app 对象。POST /api/auth/login,用户名 maintainer,密码 maintainer-password。access_token,点 Swagger 右上角 Authorize,填 Bearer 空格 token。POST /api/runs,body 使用 {"workflow":"ci_triage","prompt":"分析退款测试失败","metadata":{}}。GET /api/runs/{run_id}/trace,观察 step、model_call、tool_call、evidence。| 断点 | 你应该观察 | 继续运行后去哪里 |
|---|---|---|
app.py:create_run | request 已被 Pydantic 验证;user 已由 JWT 依赖解析。 | RunService.create |
services.py:RunService.create | fingerprint、key、created。 | repository.create_run_request |
worker.py:RunWorker.execute | delivery_id、lease 是否拿到、构造出的 RunRequest。 | AgentRuntime.run |
runtime.py:_invoke_tool | 工具名、参数、风险、超时、调用次数和返回的 ToolResult。 | EvidenceNormalizer.normalize |
先把 DEVINFRA_OFFLINE_EAGER 改为 0,创建 Run 后看它停在 pending;再改回 1 重启并重试。你会理解“API 创建任务”和“Worker 执行任务”本来是两件事。
真正的入口只有 5 行;复杂度被有意放进组合根。
"""FastAPI application entry point."""
from .composition import create_app_from_environment
app = create_app_from_environment()
devinfra_agent 包内导入”。devinfra_agent.main:app,冒号后就是它。入口只声明“我要一个 app”。composition.py 才负责根据环境选择内存实现或生产实现。这叫组合根(composition root):所有具体依赖集中装配,业务类只接收自己需要的对象。
AppSettings.from_env() 看到 DEVINFRA_MODE=offline,只解析 DEVINFRA_OFFLINE_EAGER。build_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 保存这只箱子,测试也能换入自定义假实现。
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。cls 是 AppServices 类;compose 继续把底层对象组装成领域服务。“用的到底是哪一个实现?”先看 DEVINFRA_MODE,再看 composition.build_services,最后检查 app.state.services 中对象的实际类型。
请求还没进入业务规则前,先经过验证、身份、角色、限流和统一错误处理。
FastAPI 把 JSON 解析成 LoginRequest,然后调用 container.auth.login(request)。
@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)
AppServices 对象;它在 create_app 中生成。allow(A, B) 后,函数里得到角色元组。current_user 解析 Authorization,再把 User 传进来。auth.py 的职责拆分保存用户并按用户名查找。seeded() 只为离线演示;configured() 读取生产配置。
encode 生成 header.payload.signature;decode 验签、issuer、audience、过期时间与算法。
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,写入响应,方便串联日志。/metrics 暴露文本。configure_telemetry 可安装 OpenTelemetry FastAPI instrumentor。不带 token 调一个 API 得 401;用 viewer 创建 Run 得 403;给 workflow 填不存在的值会在 Pydantic 层得 422。三个错误进入业务代码的深度不同。
同样叫“Run”,在不同边界有不同形状,这是刻意设计。
contracts.py / api_models.py
面向 API、模型网关和评测的已验证数据。非法枚举、缺字段、多字段会被拒绝。
records.py / repositories.py
Python 内部传递的数据类,如 RunRecord、ApprovalRecord,不依赖 ORM session。
models.py
SQLAlchemy 声明数据库表、列、索引、约束和 PostgreSQL JSONB/UUID 类型。
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)
uuid4 生成;list/dict 用新容器,避免多个对象共享可变默认值。StrEnum 枚举;只允许定义过的工作流字符串。
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.py 的 document_json、eval_case_json 等只做 Record → JSON 字典。Run/Event/Approval 的序列化目前在 app.py 的私有函数里。若你给数据库加字段,通常要同时检查 ORM Model、Record、迁移、Repository 转换、serializer 和 API schema。
ORM 对象依赖 session 生命周期,还可能触发隐式数据库读取。Record 把数据访问和业务层隔开,测试也能直接构造。
终于开始跟一条请求真正旅行。
@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)
RunService.create 返回已持久化的 RunRecord,不一定是刚创建的。网络重试、用户双击、代理超时都可能让客户端重复 POST。RunService.create 对影响请求含义的字段做稳定 JSON 序列化,再计算 SHA-256 fingerprint。调用方没给 key 时,使用 derived:{fingerprint};给了 key 时,数据库把 key 与 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)
如果“写数据库成功”后进程崩溃,来不及给 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
使用同一个 Idempotency-Key 连续创建相同请求,应得到同一个 Run;再保持 key 不变但修改 prompt,应得到 409 idempotency_conflict。对应测试在 test_review_durability.py。
数据层不是“存一下 JSON”,而是并发正确性的事实源。
repositories.py 有 2081 行?它提供两套相同行为:InMemoryRepository 让单测和离线模式快速确定;SqlAlchemyRepository 用真实事务、唯一约束和行锁保证跨进程正确。ThreadedRepository 再把同步 SQLAlchemy 调用放到线程,避免阻塞异步事件循环。
| 仓库能力组 | 代表方法 | 保证 |
|---|---|---|
| Run 与幂等 | create_run_request、get_run、update_run | 同 key 同请求返回旧值;同 key 异请求冲突。 |
| Outbox | claim_dispatch_outbox、mark_dispatch_* | 多 relay 竞争时一条消息只有一个 owner。 |
| 事件 | append_event、events_after | 每个 Run 内 sequence 唯一递增,可精确重放。 |
| Lease | acquire_run_lease、release_run_lease | 同一 Run 同时最多一个有效 Worker owner。 |
| 审批与动作 | claim_approval_decision、claim_action | 决定与副作用都只执行一次。 |
| Trace 投影 | create_step/model_call/tool_call/evidence/artifact/audit | projection key 去重,重投递不复制轨迹。 |
| 知识与评测 | create_document、claim_eval_run、finish_eval_run | 记录持久、评测 owner fencing。 |
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。
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"
yield 一小段文本给客户端,而不是一次性返回。Last-Event-ID,断线重连从它后面继续。进程内 EventNotifier 或生产通知机制只是让等待者早点再查一次。客户端丢包、API 重启都不怕,因为 events_after 仍能按 sequence 重放。
API 只排队,真正慢而危险的工作在 Worker。
把 ID 记入列表;eager 时用 asyncio.create_task 调 Worker。进程重启后丢失,只为本地/测试。
用稳定 task_id 给 broker 发送 devinfra.execute_run 或 eval task;Redis 是 broker/result backend。
celery_entrypoint.py 只允许 production,复用 AppSettings 和 build_services,并配置每 5 秒 relay Run outbox、每 10 秒 relay eval queue。dispatch.py 注册具体 Celery task,设置 late acknowledgement、失败重试和 delivery identity。
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)
_project_runtime 为什么存在?Runtime 只关心 Agent 逻辑,返回普通内存对象;Worker 把这些事实投影成数据库记录。每个工具原始结果先写 object_store.put,数据库只保存 URI、checksum、size,避免大日志塞进表和模型上下文。
字节存在字典,URI 使用 memory://;checksum 做内容寻址。
懒创建 bucket;同步 MinIO SDK 放到线程执行;返回 s3://bucket/key 元数据。
重投递时相同事实得到相同唯一键,Repository 把重复写变成读取旧记录。
在 runtime._invoke_tool 看原始 ToolResult;继续到 worker._project_runtime 看它变成 artifact;最后在 Trace API 看 tool_call 只剩 artifact URI/hash/size。
模型不是总指挥;Runtime 才控制顺序、预算、工具和停止条件。
nodes: Sequence[GraphNode] = (
("validate", self._validate),
("route", self._route),
("specialist", self._specialist),
("verify", self._verify),
("action_gate", self._action_gate),
("finalize", self._finalize),
)
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 --> [*]
| 对象 | 装什么 | 生命周期 |
|---|---|---|
RuntimeLimits | 最大模型调用、工具调用、replan、模型/工具/整次 Run 超时。 | Runtime 配置 |
CancellationToken | 内存取消事件;另有 durable cancellation_check 查仓库。 | 一次执行 |
_RuntimeState | request、skill、scope、证据、错误、预算计数、trace、pending action。 | 图运行/可 checkpoint |
RuntimeOutcome | 对外返回的不可变快照,供 Worker 投影。 | 图完成或中断后 |
ModelRequest/Response | 模型输入、允许工具、非信任证据、结构化输出。 | 一次模型调用 |
DeterministicGraphAdapter:普通 for 循环依次执行节点,最容易测试。LangGraphAdapter:构建 StateGraph;有 checkpointer 时支持 interrupt 与 Command(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。
这一章解释 Agent 为什么“能做事,但不能乱做事”。
把 workflow 映射到 specialist、description 和 allowed_tools。这是工具可见性的第一道上界。
名称、说明、风险、handler、超时、JSON Schema、required scope、允许角色和可选幂等 handler。
把当前用户角色、资源 scope、Skill allowlist、参数是否合法、预算是否剩余交给策略引擎。
| 风险 | 策略结果 | 例子 | 行为 |
|---|---|---|---|
| R0 | EXECUTE | repo_search、读日志 | 读操作可自动执行;瞬时错误最多 3 次有界重试。 |
| R1 | DRAFT | issue_draft | 只生成草稿,不发布真实变更。 |
| R2 | APPROVAL_REQUIRED | deployment_restart | 持久中断,人工批准后从 checkpoint 恢复;必须有下游幂等。 |
| R3 | BLOCKED | credential_export | 即使 Admin 也永久阻止自动执行。 |
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]
deployments:write。_argument_error 做项目所需的最小 JSON Schema 校验:required、额外字段、string 类型和最小长度。ToolAuditRecord;并发相同 idempotency key 共享同一 task。idempotent_handler,Fixture store 模拟下游幂等数据库。
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。
在 Run metadata 中分别指定 issue_draft、deployment_restart、credential_export,同时给对应 arguments。观察 status、decision reason、是否创建 approval 和是否出现 tool_call。
把“不可信的大段外部内容”变成“可引用、可追溯、有限大小的证据”。
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"]
KnowledgeService.create 生成 DocumentRecord,再把 KnowledgeDocument 交给 retriever.ingest。DocumentChunker 按字符上限切块并保留 overlap;每块有稳定 ID、hash、locator、metadata。tokenize/sparse_vector 形成稀疏表示。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,不把它当指令。 |
adapters/http.py
统一异步 HTTP 外壳:实际同步 urlopen 放线程,限制超时和最大响应字节。
adapters/mcp.py
校验 HTTP(S)、host allowlist、无 URL 凭据;完成 initialize、notifications/initialized、tools/list、tools/call,解析 JSON 或 SSE。
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 决策。不要第一步就怀疑模型。
API 返回 200 只说明没崩;Harness 检查 Agent 做得对不对、安全不安全。
harness/ 加载场景包,按 replay mode 执行案例、打分、聚合、比较 baseline,并由 CLI 输出 JSON/Markdown。
catalog_services.EvalService 创建 eval run 并分发;evaluation.EvalExecutionService 用 lease 执行 Harness,逐例保存 result/score,最后写 aggregate/regression。
harness/ 每个模块做什么| 模块 | 职责 | 重要对象/函数 |
|---|---|---|
| schemas.py | 案例、rubric、observation、grade、experiment 的严格 Pydantic schema。 | EvalCaseDefinition、ExperimentReport |
| dataset.py | 加载 manifest/cases/recordings,校验 72 例分布、fixture/hash 和子集。 | ScenarioPack.load/select/recording_for |
| executors.py | exact、recorded_tools、full 三种方式产生 observation。 | AgentRuntimeExecutor.execute |
| metrics.py | 精确可复现的 precision/recall/F1、Recall@K、MRR、NDCG。 | precision_recall_f1 等 |
| graders.py | 按 rubric 评分,安全项不可被 LLM judge 覆盖。 | DeterministicGrader、JudgeAugmentedGrader |
| runner.py | 循环 case/repetition,收集 CaseRunResult,聚合指标。 | HarnessRunner.run/_aggregate |
| baseline.py | 安全必须 100%,总分至少 .85,关键指标相对退化不超 .03。 | RegressionGate.evaluate |
| reporting.py / cli.py | JSON 事实报告、Markdown 展示和命令行参数。 | render_json/render_markdown/main |
直接重放最终 observation。最快、完全确定,适合 PR 门禁,但不测试当前 Runtime。
真正跑 Runtime,工具返回冻结录制值。能测规划/策略/证据链,隔离网络波动。
Runtime 与 fixture 工具全走,适合集成与鲁棒性;仍不触碰真实生产资源。
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。它让人工失败分类反哺评测集和修复优先级,但当前代码不会自动训练模型。
把离线积木逐个换成可持久、可扩展的生产积木。
进程活着就 200。容器编排据此判断是否需要重启。
关键依赖可用才 200。未 ready 时负载均衡不应发新请求。
登录后查看 repository、broker、worker、object store、Redis、MCP、Qdrant 各自状态。
ComponentCheck.required 区分关键和可降级依赖:Qdrant 可选,因为 RAG 能 lexical fallback;Redis/MCP 生产组合中是 required。HTTPHealthProbe 同样执行 scheme/host allowlist,避免健康检查变成 SSRF 工具。
| 能力 | 离线 | 生产 | 为什么换 |
|---|---|---|---|
| 业务数据 | InMemoryRepository | PostgreSQL + SQLAlchemy + Alembic | 跨进程持久、事务、唯一约束、行锁。 |
| 任务队列 | InMemoryDispatcher | Redis + Celery | 独立 Worker、重试、并发、定时 relay。 |
| 图 checkpoint | InMemorySaver | Postgres AsyncPostgresSaver | 审批等待和崩溃恢复不丢状态。 |
| 知识检索 | InMemoryHybridIndex | Qdrant + lexical fallback | 大规模 dense/sparse 向量与持久索引。 |
| 原始大对象 | InMemoryObjectStore | MinIO | 日志/diff/报告不塞数据库。 |
| 工具 | fixture registry | MCP + 可选 GitHub + fixture 补位 | 接入真实外部能力又保持统一策略。 |
| 模型 | Deterministic | Deterministic 或 OpenAI-compatible | 演示可复现;生产可选择真实模型。 |
| 限流 | 内存窗口 | Redis 计数 | 多 API 实例共享。 |
Postgres、Redis、Qdrant、MinIO、MCP 先各自健康 → API 运行 alembic upgrade head 后启动 Uvicorn → Worker 等 API 健康再启动 Celery worker + beat → Web 等 API。可观测 profile 额外启动 OTel Collector、Prometheus 和 Grafana。
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 展示。它们解决的是“发生了什么、慢在哪、请求如何串起来”,不是业务事实存储。
真实部署要用企业 IdP、Secret Manager、托管数据服务、TLS 和网络隔离。不要把固定 JWT/MCP/MinIO 密钥暴露到公网。
把“我看懂了”变成“我能安全地改”。
# 只跑一个行为,失败时显示完整信息
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
| 现象 | 先看 | 第一断点 | 相关测试 |
|---|---|---|---|
| 登录失败/401 | Authorization、token issuer/audience/exp | AuthService.authenticate | test_auth_api.py |
| 创建 Run 409 | Idempotency-Key 与 fingerprint | RunService.create | test_review_durability.py |
| Run 一直 pending | outbox 状态、dispatcher、eager/worker | OutboxRelay.relay_once | test_outbox_relay.py |
| Worker 重复/抢占 | lease owner/expiry/delivery_id | RunWorker.execute | test_delivery_leases.py |
| 工具没执行 | Skill allowlist、schema、scope、risk、budget | PolicyEngine.decide | test_policy.py |
| 审批后没恢复 | approval status/execution_state/checkpoint/action claim | _resume_approval | test_runtime_checkpoint_resume.py |
| SSE 丢/重复事件 | Last-Event-ID、sequence、terminal race | run_events.replay | test_sse.py |
| RAG 无结果 | chunk、filter、Qdrant health、degraded | HybridKnowledgeService.search | test_rag.py |
| 评测重复结果 | eval lease、projection key、owner fencing | EvalExecutionService.execute | test_evaluation_execution.py |
ToolDefinition:R0、schema、scope、role、timeout。WorkflowName 与 RUNTIME_WORKFLOWS。先找最小失败测试,再找公开调用入口,再沿对象属性进入具体实现。改完至少跑:目标测试 → 同模块测试 → 全套后端测试 → Ruff。
以后忘了某功能在哪,回到这里反查。
| 包/工具 | 项目里负责什么 | 主要出现位置 |
|---|---|---|
| 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 |
| SQLAlchemy | ORM 模型、session、查询、事务、数据库异常。 | models.py、repositories.py、composition.py |
| psycopg(间接) | SQLAlchemy 与 PostgreSQL 的驱动;也供 LangGraph saver。 | 数据库 URL |
| Alembic | 数据库 schema 版本管理与升级/降级脚本。 | backend/alembic/ |
| Celery | Run/Eval 后台任务、重试、late ack、Beat 定时 relay。 | dispatch.py、celery_entrypoint.py |
| Redis | Celery broker/result、生产限流、健康检查;不是业务事实源。 | operations.py、health.py、compose.yaml |
| LangGraph | Agent StateGraph、节点执行、interrupt/resume。 | graph.py、runtime.py |
| langgraph-checkpoint-postgres | 把图 checkpoint 持久到 PostgreSQL。 | graph.py |
| MinIO SDK | 上传/读取原始工具输出与大 artifact。 | object_store.py |
| prometheus-client | Counter/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 |
| Hatchling | 把 src/devinfra_agent 构建为 Python wheel。 | pyproject build-system |
| uv | 虚拟环境、依赖解析/锁定、运行命令;不是应用运行时库。 | uv.lock、Dockerfile、开发命令 |
asyncio 管协程/task/锁/超时;threading.RLock 保护内存仓库;collections.deque/defaultdict/Counter 管队列与统计。
dataclasses、enum.StrEnum、typing、collections.abc、uuid、datetime。
json、hashlib、hmac、base64、secrets。
urllib.request/error/parse 实现不额外依赖客户端库的 HTTP adapters;urlsplit 用于安全验证。
os 读环境变量;logging;argparse;pathlib;time;math;re。
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.error、urllib.parse、urllib.request 负责底层 HTTP 与 URL 安全解析;uuid 生成和解析实体 ID。
| 模块 | 类/函数入口(都在这里) | 职责与修改定位 |
|---|---|---|
__init__.py | 导出 main.app | 包导入便利入口;本身没有业务逻辑。 |
adapters/__init__.py | 无 | 标记 adapters 子包。 |
adapters/http.py | AdapterHttpResponse、AdapterHttpTransport.request、UrllibAdapterTransport.request/_request | 统一 HTTP transport、线程化 urlopen、超时和响应大小。 |
adapters/mcp.py | MCPStreamableHTTPClient 的 call_tool/close/tool_definitions/_initialize/_rpc/_notification/_headers/_parse_response/parse_rpc_body/_structured_from_text/_raise_http_error;_header/_bounded_redacted | MCP Streamable HTTP 生命周期、SSE/JSON 解析、认证与脱敏。 |
adapters/github.py | GitHubReadOnlyAdapter 的 get_issue/get_pull_request/get_check_runs/get_content/search_code/tool_definitions/_get/_request/_response_body/_validate_repository/_validate_number;_bounded_redacted_json/_schema | GitHub 只读工具、allowlist、缓存、路径/响应安全。 |
api_models.py | APIRequest、ToolTestRequest、DocumentCreate、KnowledgeSearchRequest、EvalCaseCreate、EvalRunCreate、EvalResultCreate、EvalScoreCreate、FailureLabelCreate | 非 Run/Approval 的 HTTP body schema。 |
app.py | AppServices.offline/compose、RunCreate、ApprovalDecision、_run_json/_event_json/_approval_json、create_app;其内 31 个 route handler | FastAPI 组合、RBAC 路由、REST/SSE 边界。新增接口先来这里。 |
auth.py | _encode/_decode/hash_password/verify_password;UserStore.seeded/by_username/configured;JWTCodec.encode/decode;AuthService.login/authenticate;auth_dependency | 本地凭据、JWT、FastAPI 当前用户依赖。 |
catalog_services.py | RegistryService.list_skills/list_tools/test_tool;KnowledgeService.create/list/search;EvalService.create_case/create_run/create_result/create_score/detail;FailureService.label/list;HealthService.components/ready | 目录、知识、评测、失败、健康的领域服务。 |
celery_entrypoint.py | 模块级 settings/services/celery_app 与 beat_schedule | Celery worker/beat 进程入口,只允许 production。 |
composition.py | AppSettings.from_env、build_services、create_app_from_environment | 环境校验、离线/生产依赖选择,排查“实际用了谁”。 |
contracts.py | Role/RiskLevel/RunStatus/WorkflowName;ContractModel/PlanStep/ToolSpec/EvidenceItem/ApprovalRequest/FinalResult/RunRequest/EvalCase/APIError/TraceEvent | API 与 Agent 共享的严格 Pydantic 合同。 |
dispatch.py | CeleryLike、CeleryDispatcher.dispatch/dispatch_eval/health、OutboxRelay.relay_once、create_celery_app、四个 register_*_task | Celery 配置、任务注册、outbox/eval relay。 |
errors.py | APIProblem、error_response、install_error_handlers 及内部 5 类 handler | 统一、安全的错误 envelope。 |
evaluation.py | EvalExecutionService.execute/_ensure_case/_config/_fail、EvalQueueRelay.relay_once | 持久 eval lease、Harness 执行、逐例结果/分数、恢复与 fencing。 |
events.py | EventPublisher、NullEventPublisher.publish、EventNotifier.publish/wait、EventService.append | 事件先持久化再通知;SSE 唤醒。 |
evidence.py | NormalizedEvidence.to_contract/model_text、EvidenceNormalizer.normalize、_safe_source_uri | citation、hash、URI、trust 和模型上下文隔离。 |
gateways.py | ToolProposal/ModelRequest/ModelResponse/ModelGateway;DeterministicModelGateway.complete;OpenAICompatibleGateway.complete;_parse_model_response/_default_transport | 模型协议、确定性假模型、OpenAI-compatible JSON Schema 调用。 |
graph.py | GraphState/GraphAdapter;DeterministicGraphAdapter.execute;LangGraphAdapter.execute/_runtime_graph/execute_runtime/resume_runtime/_execution;RuntimeGraphExecution;Postgres adapters;default_graph_adapter | 普通图、LangGraph、checkpoint interrupt/resume。 |
harness/__init__.py | 导出 Harness 常用类型 | 评测子包公共入口。 |
harness/baseline.py | RegressionGate.evaluate | 总分、安全与相对退化门禁。 |
harness/cli.py | parser、main | devinfra-eval 命令行入口。 |
harness/dataset.py | ScenarioPack.load_default/load/asset_paths/smoke_cases/robustness_cases/select/recording_for/_hash_dataset/_validate | 版本化场景包加载、hash 与完整性。 |
harness/executors.py | AgentRuntimeExecutor.execute/_evidence_ids/_recorded_registry/_full_registry | recorded/full Runtime 执行与工具注册。 |
harness/graders.py | StructuredJudge.score、DeterministicGrader.grade/_factual_score、JudgeAugmentedGrader.grade | 确定性评分与受限语义 judge。 |
harness/metrics.py | PRFScore、precision_recall_f1/tool_f1/recall_at_k/mrr/ndcg_at_k、内部 dcg | 可手算、可复现的评测指标。 |
harness/reporting.py | render_json、render_markdown | 机器事实报告与人类展示。 |
harness/runner.py | HarnessExecutor、FixtureExecutor.execute、HarnessRunner.run/_aggregate | 案例/重复循环、评分与聚合。 |
harness/schemas.py | HarnessModel、ReplayMode/EvalWorkflow、GoldEvidence/EvalRubric/EvalCaseDefinition/ScenarioManifest/CaseObservation/CaseGrade/ExperimentConfig/Metadata/CaseRunResult/Report/RegressionDecision | 全部评测数据合同与交叉字段校验。 |
health.py | ComponentCheck、HTTPHealthProbe.health、RedisHealthProbe.health | 外部依赖健康探测与 SSRF 边界。 |
main.py | app = create_app_from_environment() | Uvicorn API 入口。 |
models.py | Base、两个 mixin、21 个 ORM Model | 数据库表、列、索引、唯一约束和 PostgreSQL 类型。 |
object_store.py | StoredObject、_validate_key、InMemoryObjectStore.put/get/health、MinioObjectStore._ensure_bucket/put/get/health | 原始 artifact 内容寻址与对象存储。 |
operations.py | JsonLogFormatter.format、configure_structured_logging、两种 RateLimiter、HttpMetrics.observe/render、install_operations、configure_telemetry | 日志、request ID、限流、Prometheus、OTel。 |
policy.py | ActionDecision、PolicyContext、PolicyDecision、PolicyEngine.decide | fail-closed R0–R3 决策。 |
rag.py | KnowledgeDocument/Chunk/SearchHit/Response、Embedding/Index protocols、两种 embedding、DocumentChunker.chunk、InMemoryHybridIndex、HybridKnowledgeService.ingest/search、QdrantHTTPStore、tokenize/sparse_vector/bm25_scores/rank_scores/cosine/injection_suspected/chunk_payload/qdrant_hit | 完整混合检索、降级与引用。 |
records.py | utc_now;Document/Eval/Feedback/RunStep/ModelCall/ToolCall/Evidence/Artifact/Audit/Skill/Tool/Chunk 共 16 个 Record | 跨层内部数据载体。 |
repositories.py | _aware;Run/Event/Approval/Outbox/ActionClaim records;InMemoryRepository 与 SqlAlchemyRepository 的 Run、outbox、cancel、event、lease、approval、action、projection、catalog、eval 全套方法;全部 _xxx_record/_xxx_model 转换;ThreadedRepository;new_run | 业务事实与并发正确性核心。详细能力组见第 08 章。 |
serializers.py | document_json/eval_case_json/eval_run_json/eval_result_json/eval_score_json/feedback_json | Record 到 API JSON。 |
services.py | Dispatcher、InMemoryDispatcher.dispatched/dispatch/dispatch_eval/health、scopes_for_role、RunService.create/rerun/cancel/require、ApprovalService.create_for_run/decide/validate_action | Run 生命周期、幂等创建、取消、审批业务规则。 |
skills.py | RUNTIME_WORKFLOWS、SkillConfig.from_mapping、DEFAULT_SKILL_CONFIGS、SkillRegistry.default/get/workflows | workflow → specialist + 最小工具集。 |
tools.py | ToolRisk/InputError/TransientError/Result/AuditRecord、两种 Handler 类型、FixtureDownstreamIdempotencyStore.execute/effect_count/result、ToolDefinition、ToolRegistry.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.py | RuntimeLimits、CancellationToken、RunCancelled、RuntimeOutcome、_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_operation | Agent 的真正控制器:状态、图、预算、模型、工具、验证、审批、取消。 |
worker.py | RunLeaseBusy、RunWorker.cancel/health/execute/_resume_approval/_project_runtime | lease、Runtime 执行、审批恢复、Trace/artifact 投影。 |
上表按“先找功能”压缩了重复方法。下面把剩余函数、内部函数、类型和状态变量补齐;这些名称也可被页面搜索找到。
| 位置 | 补充符号 | 统一解释 |
|---|---|---|
| dispatch.py | CeleryLike.send_task;register_run_task 内的 execute_run;register_outbox_task 内的 relay_outbox;register_eval_task 内的 execute_eval;register_eval_relay_task 内的 relay_eval_queue | send_task 是 Celery 最小协议;四个 register 函数把普通 Python worker/relay 包成 Celery task,内部函数是 Celery 真正调用的入口。 |
| errors.py | api_problem_handler、validation_handler、http_handler、integrity_handler、unexpected_handler | 分别捕获项目业务错误、Pydantic/FastAPI 校验错、HTTPException、SQLAlchemy 完整性错误和兜底异常,最终都调用统一 error response。 |
| gateways.py / graph.py | _default_transport.send;GraphState.check_cancelled;AsyncPostgresLangGraphAdapter._runtime_call | send 是送进线程的同步 urlopen 闭包;check_cancelled 是图状态协议要求;_runtime_call 统一打开 saver、setup、转调 execute/resume 并关闭资源。 |
| harness/executors.py | transient_handler、timeout_handler、malicious_runbook | 三者是构造评测工具的内部 handler,分别模拟瞬时失败、永不及时返回和携带 prompt injection 的 Runbook,用来验证重试、超时和安全隔离。 |
| harness/schemas.py | EvalRubric.validate_weight_sum、EvalCaseDefinition.validate_tool_sets、ExperimentMetadata | 两个 Pydantic model validator 做跨字段校验:权重和必须合法,expected/forbidden tool 不能矛盾;ExperimentMetadata 记录数据集 hash、模式、时间等实验身份。 |
| operations.py | InMemoryRateLimiter、RedisRateLimiter、prometheus_metrics、operations_middleware | 两种 limiter 实现同一窗口语义;metrics handler 返回 Prometheus bytes;middleware 包住每个 HTTP 请求,做 request ID、限流、计时和统一响应头。 |
| rag.py | VectorStoreUnavailable;KnowledgeChunk、SearchResponse;EmbeddingProvider、DeterministicEmbeddingProvider、MultilingualMiniLMEmbeddingProvider;upsert、upsert_documents、search_lexical、_filtered、ensure_collection | 异常触发显式降级;Chunk/Response 是检索数据;Provider 协议允许假/真 embedding 互换;upsert 系列写索引,lexical/filter 做内存退路,ensure_collection 幂等创建 Qdrant collection。 |
| models.py | IdentityMixin、TimestampMixin;UserModel、RunModel、RunEventModel、RunStepModel、ModelCallModel、ToolCallModel、EvidenceModel、ArtifactModel、ApprovalModel、AuditLogModel、SkillModel、ToolModel、DocumentModel、ChunkModel、EvalCaseModel、EvalRunModel、EvalResultModel、EvalScoreModel、FeedbackModel、DispatchOutboxModel、ActionClaimModel | 两个 mixin 复用主键与时间列;其余每个 Model 一一对应数据库表。名字去掉 Model 后就是业务实体,字段细节以类内 mapped_column 为准。 |
| records.py / repositories.py | EvalCaseRecord、EvalRunRecord、EvalResultRecord、EvalScoreRecord、FeedbackRecord、RunStepRecord、ModelCallRecord、ToolCallRecord、EvidenceRecord、ArtifactRecord、SkillRecord、ToolRecord、ChunkRecord;EventRecord、DispatchOutboxRecord、ActionClaimRecord | Record 是离开 ORM session 后仍可安全使用的普通数据对象;名称与 Model 对应。Event/Outbox/ActionClaim 定义在 repositories.py,因为它们紧贴持久并发语义。 |
| repositories.py:Run/Outbox | list_dispatch_outbox、_queue_dispatch_locked/_queue_dispatch、request_cancel、is_cancel_requested、record_worker_failure、get_run_by_idempotency_key、run_lease_retry_at | 列出/排队 outbox;持久取消;失败计数并在上限后终止;按幂等 key 找 winner;告诉 Celery 何时能在 lease 过期后重试。带 locked 的内存版本要求调用者已持有 RLock。 |
| repositories.py:审批/动作 | create_approval、approval_for_run、begin_approval_execution、complete_approval_execution、mark_approval_checkpoint_resumed、complete_action、get_action_claim | 从创建审批到 claim decision、claim execution、标记 checkpoint 已恢复、写动作 receipt 的状态机。每一步都检查旧状态与 owner,防止重复副作用。 |
| repositories.py:评测/反馈 | get_eval_case、get_eval_run、get_eval_result、list_eval_results、list_eval_scores、create_feedback、list_feedback | 按主键读评测实体、按父 ID 列子项,保存/读取人工失败标签。 |
| repositories.py:Trace/Catalog | set_health;list_steps、create_model_call/list_model_calls、create_tool_call/list_tool_calls、create_evidence/list_evidence、create_artifact/list_artifacts、create_audit/list_audits、upsert_skill/list_skill_records、upsert_tool/list_tool_records、create_chunk/list_chunks | set_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_record;ThreadedRepository.__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_checkpoint;AgentRuntime.restore_checkpoint_state/awaiting_approval_status/finalize_checkpoint_state | __post_init__ 验证预算/超时为正;cancelled 是只读状态;State 方法负责内存+持久取消、记录 partial error 和可序列化 checkpoint;AgentRuntime 三个桥接方法供 graph.py 在恢复时调用。 |
| tools.py | ToolInputError、TransientToolError | 前者表示参数/幂等前置条件错误,不应盲重试;后者明确表示 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/恢复字段与唯一性。 | 历史迁移,不应直接改。 |
| 方法与路径 | 处理函数 | 进入的服务/行为 |
|---|---|---|
POST /api/auth/login | login | AuthService.login,返回 Bearer JWT。 |
GET /api/auth/me | me | 返回当前依赖注入的 User 身份。 |
POST /api/runs | create_run | RunService.create,幂等持久化并分发。 |
GET /api/runs | list_runs | repository.list_runs。 |
GET /api/runs/{run_id} | run_detail | RunService.require。 |
GET /api/runs/{run_id}/trace | run_trace | 并发读取 events、steps、model/tool calls、evidence、artifacts。 |
GET /api/runs/{run_id}/events | run_events | 按 Last-Event-ID 重放的 SSE generator。 |
POST /api/runs/{run_id}/rerun | rerun | RunService.rerun,保留原始输入并记录 rerun_of。 |
POST /api/runs/{run_id}/cancel | cancel | RunService.cancel,内存+持久取消。 |
GET /api/approvals | list_approvals | repository.list_approvals。 |
GET /api/approvals/{approval_id} | approval_detail | repository.get_approval,不存在转 404。 |
POST /api/approvals/{approval_id}/decision | decide_approval | ApprovalService.decide,只允许 Maintainer/Admin。 |
GET /api/skills | list_skills | RegistryService.list_skills。 |
GET /api/tools | list_tools | RegistryService.list_tools。 |
POST /api/tools/{name}/test | test_tool | RegistryService.test_tool,仍经过 Policy。 |
POST /api/knowledge/documents | create_document | KnowledgeService.create,保存并 ingest。 |
GET /api/knowledge/documents | list_documents | KnowledgeService.list。 |
POST /api/knowledge/search | search_knowledge | KnowledgeService.search,返回 hits/degraded/limitations。 |
POST /api/evals/cases | create_eval_case | EvalService.create_case。 |
GET /api/evals/cases | list_eval_cases | repository.list_eval_cases。 |
POST /api/evals/runs | create_eval_run | EvalService.create_run,持久化后分发。 |
GET /api/evals/runs | list_eval_runs | repository.list_eval_runs。 |
GET /api/evals/runs/{eval_run_id} | eval_run_detail | EvalService.detail,嵌套 results/scores。 |
POST /api/evals/runs/{eval_run_id}/results | create_eval_result | EvalService.create_result。 |
POST /api/evals/results/{eval_result_id}/scores | create_eval_score | EvalService.create_score。 |
GET /api/failures | list_failures | FailureService.list,附人工 labels。 |
POST /api/failures/{run_id}/label | label_failure | FailureService.label。 |
GET /api/health/components | component_health | HealthService.components,需要登录。 |
GET /healthz | healthz | 无依赖的 liveness。 |
GET /readyz | readyz | HealthService.ready,失败返回 503。 |
GET /health | legacy_health | 兼容旧调用方的健康接口。 |
今后先从现象确定领域,再从 API/Worker/Harness 入口进入,用本表找类和方法。不要从文件数量最大的模块开始盲读。