实战 02:简易 OpenCode——Python 编码 Agent
1、本篇交付物
使用 Python、FastAPI、SQLite、Pydantic 和 LangGraph 完成一个最小 Server:创建 Session、发送 Prompt、SSE 订阅事件、执行只读工具、产生编辑权限请求、批准后恢复并应用 Patch。真实模型最后接入,开发阶段先用 FakeModel 保证测试确定。
2、项目结构
opencode_py/
app.py # FastAPI 路由和生命周期
contracts.py # Session/Message/Part/Event schema
repository.py # SQLite 事务
event_bus.py # 持久事件 + 本地订阅
providers/base.py # 模型统一接口
providers/fake.py
context_builder.py # 规则、消息、压缩和 token 预算
tools/base.py
tools/files.py
permissions.py
workspace.py
graph.py
tests/
3、核心合同
from datetime import datetime
from typing import Annotated, Literal, Union
from pydantic import BaseModel, Field
class Session(BaseModel):
id: str
project_id: str
mode: Literal["plan", "build"] = "plan"
status: Literal["idle", "running", "waiting_permission", "failed"] = "idle"
last_event_seq: int = 0
class TextPart(BaseModel):
type: Literal["text"] = "text"
id: str
text: str = ""
class ToolPart(BaseModel):
type: Literal["tool"] = "tool"
id: str
call_id: str
tool: str
status: Literal["pending", "waiting_permission", "running", "completed", "error"]
input: dict
output_artifact_id: str | None = None
Part = Annotated[Union[TextPart, ToolPart], Field(discriminator="type")]
class Event(BaseModel):
session_id: str
seq: int
type: str
payload: dict
created_at: datetime
API、数据库 JSON 和 SSE 都使用这些 schema。数据库读取出的 JSON 仍要 model_validate,不能因为是自己的数据就跳过版本/结构检查。
4、SQLite 事务和事件序号
每次状态变更与事件写入同一事务:
def append_event(conn, session_id: str, event_type: str, payload: dict) -> Event:
row = conn.execute(
"UPDATE sessions SET last_event_seq=last_event_seq+1 WHERE id=? RETURNING last_event_seq",
(session_id,),
).fetchone()
if row is None:
raise LookupError("SESSION_NOT_FOUND")
event = Event(session_id=session_id, seq=row[0], type=event_type,
payload=payload, created_at=datetime.utcnow())
conn.execute(
"INSERT INTO events(session_id,seq,type,payload,created_at) VALUES(?,?,?,?,?)",
(session_id, event.seq, event.type, event.model_dump_json(), event.created_at.isoformat()),
)
return event
单 Session 序号严格递增;UNIQUE(session_id, seq) 防止并发重复。事务提交后再通知内存 subscriber;订阅者漏掉通知也能从 events 表按 seq 补读。
5、SSE:先补历史,再等待新事件
from fastapi.responses import StreamingResponse
def encode_sse(event: Event) -> str:
return f"id: {event.seq}\nevent: {event.type}\ndata: {event.model_dump_json()}\n\n"
@app.get("/sessions/{session_id}/events")
async def events(session_id: str, after: int = 0, actor=Depends(auth)):
require_session_owner(actor, session_id)
async def stream():
cursor = after
while True:
batch = repo.events_after(session_id, cursor, limit=100)
for event in batch:
cursor = event.seq
yield encode_sse(event)
if batch:
continue
await event_bus.wait(session_id, cursor, timeout=15)
yield ": keepalive\n\n"
return StreamingResponse(stream(), media_type="text/event-stream")
生产要处理客户端取消、subscriber 清理、最大事件保留和窗口过期。SSE 不携带密钥与内部 stack。
6、模型 Provider 与 FakeModel
class ModelChunk(BaseModel):
type: Literal["text_delta", "tool_call", "finish"]
text: str | None = None
call_id: str | None = None
tool: str | None = None
arguments: dict | None = None
reason: str | None = None
class ModelProvider(Protocol):
async def stream(self, messages: list[dict], tools: list[dict], *, signal) -> AsyncIterator[ModelChunk]: ...
FakeModel 根据测试脚本依次返回 text/tool call/finish,可稳定覆盖审批、重连和错误。真实 Provider adapter 负责把 LangChain/provider chunk 转为内部 ModelChunk,并记录 usage、模型版本和原始 request ID。
7、Tool Registry
class ToolContext(BaseModel):
session_id: str
project_root: str
workspace_root: str
mode: Literal["plan", "build"]
class ToolSpec(BaseModel):
name: str
description: str
input_schema: dict
permission: Literal["read", "edit", "execute", "network"]
class Tool(Protocol):
spec: ToolSpec
async def execute(self, raw: dict, ctx: ToolContext) -> dict: ...
read_file 输入只有相对 path、offset、limit;root 从 ToolContext 获得。先做路径规范化、symlink/保留目录/大小检查,再读取。grep 使用 argv 调 rg 或受控纯 Python 实现,并限制 pattern、glob、结果数和超时。
apply_patch 不直接相信 patch 文本:解析每个 marker 路径、验证工作区、文件 hash、总修改量并 dry-run。写入前生成 snapshot 和 Patch artifact。
8、Permission Engine
class Rule(BaseModel):
tool: str
pattern: str = "*"
action: Literal["allow", "ask", "deny"]
scope: Literal["global", "project", "session"]
def decide(rules: list[Rule], tool: ToolSpec, input_summary: str, mode: str) -> str:
if mode == "plan" and tool.permission in {"edit", "execute"}:
return "deny"
matches = [r for r in rules if fnmatch(tool.name, r.tool) and fnmatch(input_summary, r.pattern)]
matches.sort(key=specificity, reverse=True)
return matches[0].action if matches else "ask"
ask 时在事务中创建 permission request、更新 ToolPart 和 Session 状态,然后图 interrupt。响应 API 验证 owner、pending 状态、call/input hash 和过期时间;“记住”只创建同等或更窄规则。
9、LangGraph Agent Loop
class AgentState(TypedDict, total=False):
session_id: str
turn_id: str
messages: list[dict]
pending_call: dict
permission_request_id: str
budget: dict
done: bool
def call_model(state, runtime): ... # 流式写 Parts/Events,返回 pending_call 或 done
def check_permission(state, runtime): ... # allow/deny/interrupt
def run_tool(state, runtime): ... # 写 ToolPart 和 artifact
def route(state): return END if state["done"] else "check_permission" if state.get("pending_call") else "call_model"
builder = StateGraph(AgentState)
builder.add_node("call_model", call_model)
builder.add_node("check_permission", check_permission)
builder.add_node("run_tool", run_tool)
builder.add_edge(START, "call_model")
builder.add_conditional_edges("call_model", route)
builder.add_conditional_edges("check_permission", permission_route)
builder.add_edge("run_tool", "call_model")
graph = builder.compile(checkpointer=production_checkpointer)
数据库 Session/Message 是产品权威;Graph state 保存当前 turn 的执行游标。节点重试必须通过 call ID/tool execution 唯一约束避免重复写。
10、API
POST /projects/open 验证并注册本地项目
POST /sessions 创建 Session
GET /sessions/:id snapshot:Session + Messages + Parts + pending permission
POST /sessions/:id/messages 入库用户消息并启动后台 turn
GET /sessions/:id/events?after=n SSE
POST /permissions/:id/respond allow_once / allow_rule / deny
POST /sessions/:id/cancel 取消当前 turn
POST /sessions/:id/undo 应用反向 Patch
发送消息接口返回 202,不等待模型;每个 Session 同时只允许一个活跃 turn,或明确实现排队。
11、测试顺序
先测试 repository 事务和 seq;再测试 safe path、grep、Patch;再用 FakeModel 跑完整图;最后启动 FastAPI 做 SSE/权限集成测试。关键场景:客户端断线不重跑;Plan 模式 edit 被拒绝;批准后文件已变化返回 stale;相同 call ID 不重复执行;取消能结束模型与命令;Server 重启后从 SQLite/checkpoint 恢复。
做到这一步,Python 版已经是一个可运行的本地编码 Agent Server。后续 TUI、上下文压缩、LSP/MCP 和 Sandbox 在专题篇继续完善。
如果您觉得这篇文章有帮助,请点个赞吧~
评论
请登录后发表评论
去登录