openai-agents-python 沙箱会话事件 Sinks 全解析:从回调、JSONL 落盘到 HTTP 代理的事件投递体系
【免费下载链接】openai-agents-pythonA lightweight, powerful framework for multi-agent workflows项目地址: https://gitcode.com/GitHub_Trending/op/openai-agents-python
本篇文章聚焦 openai-agents-python 沙箱子系统中负责会话审计事件投递的 Sinks 机制(sinks 模块)。通过阅读本文,你将掌握EventSink抽象、五种内置 Sink(回调、JSONL 文件、工作区 JSONL、HTTP 代理、链式组合)的完整参数与适用场景,理解DeliveryMode、OnErrorPolicy、EventPayloadPolicy等配套机制如何协同控制事件的投递方式、容错策略与敏感数据裁剪,并学会如何将沙箱中的每次exec、write等操作接入自己的日志、审计与监控管道。
一、Sinks 在沙箱体系中的定位
openai-agents-python 的沙箱(Sandbox)模块为 Agent 提供了隔离的代码执行与文件系统操作能力。当 Agent 在沙箱会话(SandboxSession)中执行命令、读写文件时,每一次操作都会产生结构化的审计事件(SandboxSessionEvent)。Sinks 的角色,就是这些事件的下游消费者(consumer)——它们决定事件最终被写进日志文件、转发到远程服务,还是直接交给用户提供的回调函数。
从模块依赖看,sinks.py 建立在两个关键基座上:
- 事件模型events.py:定义了
SandboxSessionStartEvent/SandboxSessionFinishEvent等事件结构,两者通过phase字段构成判别联合(discriminated union); - 会话抽象base_sandbox_session.py:提供
write()、read()、running()、register_persist_workspace_skip_path()等底层原语,工作区型 Sink(WorkspaceJsonlSink)正是通过这些原语把事件写进沙箱内部。
Sinks 的投递由 manager.py 中的Instrumentation类驱动:它持有 Sink 列表,在每次操作开始与结束时调用emit(),将事件按顺序分发给每个 Sink。同时,SandboxSession包装器(sandbox_session.py)在构造时会调用_bind_session_to_sinks(),把实现了SandboxSessionBoundSink协议的 Sink 绑定到底层会话上,从而让 Sink 能访问会话状态(如session_id)。
二、核心抽象:EventSink 与配套协议
1.EventSink抽象基类
所有 Sink 都继承自 EventSink,其接口非常精简:
class EventSink(abc.ABC): """Consumes SandboxSessionEvent objects (e.g., callback, file outbox, proxy HTTP).""" name: str | None = None mode: DeliveryMode on_error: OnErrorPolicy payload_policy: EventPayloadPolicy | None @abc.abstractmethod async def handle(self, event: SandboxSessionEvent) -> None: ...每个 Sink 需要声明四个属性:
| 属性 | 类型 | 含义 |
|---|---|---|
name | str \| None | Sink 的可选名称,便于日志与调试识别 |
mode | DeliveryMode | 投递模式,见下文 |
on_error | OnErrorPolicy | 出错时的处理策略 |
payload_policy | EventPayloadPolicy \| None | 事件负载裁剪策略,控制敏感/大体积数据是否进入事件 |
handle()是唯一的抽象方法,负责把单个事件投递到目标。
2. 关键类型别名
模块顶部定义了两个受控的字面量类型(sinks.py):
DeliveryMode = Literal["sync", "async", "best_effort"] OnErrorPolicy = Literal["raise", "log", "ignore"]DeliveryMode(投递模式):sync:同步投递,调用方等待 Sink 处理完当前事件;async:异步投递,事件处理在后台任务中进行,不阻塞主流程;best_effort:尽力而为,投递失败不影响会话主流程,通常配合log容错策略使用。
OnErrorPolicy(错误策略):raise(向上抛出异常)、log(记录日志后继续)、ignore(静默忽略)。
3.SandboxSessionBoundSink协议
部分 Sink 需要访问底层会话,例如读取session_id或调用write()写入工作区。为此模块定义了运行时可检查的协议(sinks.py):
@runtime_checkable class SandboxSessionBoundSink(Protocol): """Optional interface for sinks that need access to the underlying SandboxSession.""" def bind(self, session: BaseSandboxSession) -> None: ...SandboxSession在构造时遍历Instrumentation.sinks,对实现该协议的 Sink(包括ChainedSink内部的子 Sink)调用bind(self._inner),绑定到底层会话而非包装器,从而避免递归的事件循环(详见 sandbox_session.py)。
4._unwrap_session_wrapper防御性解包
sinks.py 中还有一个防御性工具函数:如果 Sink 被意外绑定到了SandboxSession包装器上,它会解包到内部的BaseSandboxSession,防止事件在包装层反复触发产生递归循环。该函数通过检查类名与模块名来识别包装器,避免引入循环依赖。
三、五种内置 Sink 详解
1.CallbackSink:把事件交给用户回调
CallbackSink 将事件直接投递给用户提供的可调用对象,同时支持同步与异步回调(检测到协程对象时会await):
CallbackSink( callback: Callable[[SandboxSessionEvent, BaseSandboxSession], object], *, mode: DeliveryMode = "sync", on_error: OnErrorPolicy = "raise", payload_policy: EventPayloadPolicy | None = None, name: str | None = None, )关键行为:
- 回调签名接收两个参数:事件对象和已绑定的会话对象,方便在回调中读取会话上下文;
- 未绑定会话时调用
handle()会抛出RuntimeError,提示需要使用带插桩的SandboxSession或调用bind(session); - 适用场景:内存中的事件收集器、实时告警、与自定义观测系统对接。
测试用例 test_session_sinks.py 展示了典型用法——构造Instrumentation时传入CallbackSink(lambda e, _sess: events.append(e), mode="sync"),随后在async with SandboxSession(...)中执行exec("echo hi"),即可从events列表中取回op == "exec"的完成事件并断言其stdout内容。
2.JsonlOutboxSink:追加写入宿主机 JSONL 文件
JsonlOutboxSink 把每个事件序列化为一行 JSON,追加到宿主机文件系统上的指定文件:
JsonlOutboxSink( path: Path, *, mode: DeliveryMode = "best_effort", on_error: OnErrorPolicy = "log", payload_policy: EventPayloadPolicy | None = None, )实现细节:
- 使用
event_to_json_line()(utils.py)生成紧凑 JSON 行:json.dumps(..., separators=(",", ":"), sort_keys=True)并附加换行符,保证每行独立可解析; - 写入发生在线程池中(
asyncio.to_thread),避免阻塞事件循环; - 自动创建父目录(
mkdir(parents=True, exist_ok=True)); - 在 POSIX 平台会尝试用
fcntl.flock对文件加排他锁,防止多进程并发追加时互相覆盖,非 POSIX 平台(如 Windows)自动降级为不加锁; - 默认
best_effort+log的组合意味着日志写入失败不会拖垮沙箱主流程。
测试 test_jsonl_outbox_sink_appends_one_line_per_event 验证了"每个事件恰好追加一行"的契约,并通过ChainedSink与回调 Sink 组合验证顺序投递。
3.WorkspaceJsonlSink:写入沙箱工作区内部
WorkspaceJsonlSink 与前者最大的区别是:事件日志不落在宿主机,而是写进会话工作区(manifest.root之下)。它仍在客户端进程运行,但通过SandboxSession.write()写入会话,因此对 Docker / Modal 等无宿主机挂载卷的沙箱同样适用。
WorkspaceJsonlSink( *, workspace_relpath: Path = Path("logs/events-{session_id}.jsonl"), ephemeral: bool = False, mode: DeliveryMode = "best_effort", on_error: OnErrorPolicy = "log", payload_policy: EventPayloadPolicy | None = None, flush_every: int = 1, )参数说明:
| 参数 | 默认值 | 说明 |
|---|---|---|
workspace_relpath | logs/events-{session_id}.jsonl | 工作区根目录下的相对路径,支持轻量模板,在bind()时展开 |
ephemeral | False | 为True时,通过register_persist_workspace_skip_path()将日志路径排除出未来的工作区快照,适合临时审计输出 |
flush_every | 1 | 每积累 N 个事件触发一次落盘刷新,最小为 1 |
路径模板:workspace_relpath支持两个占位符,在bind()时用会话状态展开(sinks.py):
{session_id}:UUID 字符串,如550e8400-e29b-41d4-a716-446655440000;{session_id_hex}:无连字符的十六进制 UUID,如550e8400e29b41d4a716446655440000。
示例:Path("logs/events-{session_id}.jsonl")会渲染为logs/events-550e8400-e29b-41d4-a716-446655440000.jsonl。
缓冲与刷盘策略:handle()内部先把事件编码进内存缓冲(bytearray),仅在满足条件时才真正写盘(sinks.py):
- 缓冲条数达到
flush_every的整数倍; - 遇到
persist_workspace的start阶段; - 遇到
stop操作; - 遇到
shutdown的start阶段(shutdown的finish阶段会抑制本次刷新,因为此时底层沙箱可能已不可写)。
时序安全:写盘前会通过_can_flush_to_workspace()检查session.running()。这是因为SandboxSession.start()在底层沙箱完全就绪前就会发出start事件,早期启动或收尾阶段写入可能失败(sinks.py)。每次刷新采用"读取已有内容 + 追加新内容 + 整体写回"的方式,读取时对FileNotFoundError/WorkspaceReadNotFoundError做了容错,视为空文件处理。
测试 test_workspace_jsonl_sink_writes_into_workspace_and_persists 验证了默认路径下日志文件能被正确写入并持久化。
4.HttpProxySink:POST 事件到代理端点
HttpProxySink 将每个事件以 JSON 形式 POST 到指定的代理端点(本地守护进程或远程服务):
HttpProxySink( endpoint: str, *, headers: dict[str, str] | None = None, timeout_s: float = 5.0, spool_path: Path | None = None, mode: DeliveryMode = "best_effort", on_error: OnErrorPolicy = "log", payload_policy: EventPayloadPolicy | None = None, )关键行为:
- 请求体为
event.model_dump_json().encode("utf-8"),并自动附带content-type: application/json头,可与自定义headers合并; timeout_s控制单次请求超时(默认 5 秒),网络调用在后台线程执行;- 失败溢出(spool):当
spool_path提供且 POST 失败(OSError)时,事件行会先被追加写入本地的 spool 文件(同样自动建父目录、逐行 flush),随后抛出RuntimeError提示 "http proxy sink POST failed",便于上层按on_error策略处理; - 典型场景:把沙箱审计事件转发给本地收集代理、可观测性网关或审计服务。
5.ChainedSink:按序组合多个 Sink
ChainedSink 用于把多个 Sink 编成一组、按顺序执行:
ChainedSink(*sinks: EventSink)设计要点:
- 直接调用时,
handle()会依次await每个子 Sink 的handle(),保证顺序; - 更常见的是由
Instrumentation解包使用:emit()检测到ChainedSink后展开其内部 Sink,为每个子 Sink单独应用 per-op/per-sink 负载策略,再保证顺序投递——即"分组不会禁用每个 Sink 的策略行为"(manager.py); - 为保持
EventSink接口一致性,其自身的mode/on_error/payload_policy被固定为占位值("sync"/"raise"/None),实际行为由解包路径决定。
四、事件模型与负载策略
1. 事件结构
所有事件共享 SandboxSessionEventBase 的公共字段:
| 字段 | 类型 | 说明 |
|---|---|---|
version | int | 事件格式版本,默认1 |
event_id | uuid.UUID | 事件唯一 ID |
ts | datetime | UTC 时间戳 |
session_id | uuid.UUID | 所属会话 ID |
seq | int | 会话内事件序号,用于还原顺序 |
op | OpName | 操作名(如exec、write) |
phase | "start" \| "finish" | 阶段 |
span_id/parent_span_id/trace_id | str \| None | 与 SDK 追踪联动的 span / trace 标识 |
data | dict | 操作相关的元数据(路径、argv、耗时等) |
SandboxSessionFinishEvent额外携带ok、duration_ms、error_code、error_type、error_message、error_retryable以及可选的stdout/stderr字符串;原始字节输出(stdout_bytes/stderr_bytes)默认被exclude=True排除在序列化之外,仅用于 per-sink/per-op 策略裁剪(events.py)。
事件通过phase判别联合统一为SandboxSessionEvent,可用validate_sandbox_session_event()从任意对象解析出正确的阶段模型(events.py)。
2.EventPayloadPolicy:敏感数据裁剪
EventPayloadPolicy 控制事件中包含多少潜在敏感或大体积数据:
class EventPayloadPolicy(BaseModel): include_exec_output: bool = Field(default=False) # exec 输出默认关闭 max_stdout_chars: int = Field(default=8_000, ge=0) max_stderr_chars: int = Field(default=8_000, ge=0) include_write_len: bool = Field(default=True) # write 事件仅含尽力而为的字节数- exec 输出默认关闭:命令输出可能既嘈杂又敏感,需要显式开启
include_exec_output=True才会把截断后的stdout/stderr放进事件; - 截断:
_safe_decode()以 UTF-8 容错解码(errors="replace"),按解码后字符串长度截断到max_stdout_chars/max_stderr_chars,超出部分以…结尾,保证 JSON 始终合法(utils.py); - write 事件不含文件字节:只尽力而为地统计写入长度(通过 seekable 流的
tell/seek估算,见 utils.py),绝不把文件内容放进审计事件。
Instrumentation支持全局策略payload_policy与按操作细分的payload_policy_by_op(dict[OpName, EventPayloadPolicy]),后者可对特定操作覆盖全局配置(manager.py)。
测试 test_sandbox_session_exec_emits_stdout_when_enabled 展示了开启include_exec_output=True后exec("echo hi")的完成事件能取回 stdout 且span_id以sandbox_op_开头。
五、Instrumentation:事件分发的总控
Instrumentation 是 Sink 机制的调度中枢:
Instrumentation( *, sinks: Sequence[EventSink] | None = None, payload_policy: EventPayloadPolicy | None = None, payload_policy_by_op: dict[OpName, EventPayloadPolicy] | None = None, )- 持有 Sink 列表(
sinks属性返回副本),并提供add_sink()动态追加; emit()按顺序遍历 Sink 投递事件:遇到ChainedSink时解包,为每个内层 Sink 计算策略、应用裁剪后再保证顺序投递;- 异步投递的任务被记录在
_tasks集合中,便于跟踪未完成任务。
SandboxSession的instrumented_op()装饰器负责在每次操作的前后构造SandboxSessionStartEvent/SandboxSessionFinishEvent并调用Instrumentation.emit(),同时把 SDK 追踪的 span/trace 信息写入事件字段,实现审计事件与追踪系统的联动(sandbox_session.py)。
六、组合实战:一套完整的沙箱审计管道
综合以上组件,可以搭建"宿主机 JSONL 归档 + 工作区内审计日志 + 实时回调监控"的管道:
import asyncio from pathlib import Path from agents.sandbox.session import ( CallbackSink, ChainedSink, Instrumentation, JsonlOutboxSink, WorkspaceJsonlSink, EventPayloadPolicy, ) from agents.sandbox.sandboxes.unix_local import ( UnixLocalSandboxSession, UnixLocalSandboxSessionState, ) def on_event(event, session) -> None: # 实时监控:例如 op 失败时打印告警 if event.phase == "finish" and not getattr(event, "ok", True): print(f"[audit] op={event.op} failed: {event.error_message}") instrumentation = Instrumentation( sinks=[ ChainedSink( # 宿主机侧 JSONL 归档(尽力而为) JsonlOutboxSink(Path("/tmp/sandbox-events.jsonl")), # 工作区内部审计日志,排除出快照 WorkspaceJsonlSink( workspace_relpath=Path("logs/events-{session_id}.jsonl"), ephemeral=True, ), # 同步回调,事件失败实时可见 CallbackSink(on_event, mode="sync", on_error="raise"), ) ], # 开启 exec 输出(截断到默认 8000 字符) payload_policy=EventPayloadPolicy(include_exec_output=True), ) # 以 UnixLocal 沙箱为例(Docker 沙箱用法同理) inner = UnixLocalSandboxSession.from_state( UnixLocalSandboxSessionState(manifest=Manifest(root="/tmp/ws")) ) async with SandboxSession(inner, instrumentation=instrumentation) as session: await session.exec("echo hello from sandbox") await session.write(Path("note.txt"), io.BytesIO(b"data"))执行要点:
ChainedSink保证三个 Sink 按声明顺序依次收到每个事件;CallbackSink使用mode="sync"保证监控实时性,JsonlOutboxSink/WorkspaceJsonlSink默认best_effort,日志管道故障不会中断沙箱主流程;ephemeral=True让工作区内的审计日志不进入未来快照,避免审计数据污染持久化工作区(底层依赖 base_sandbox_session.py 的register_persist_workspace_skip_path);- 所有事件类、Sink 类均已从
agents.sandbox.session顶层导出(见init.py)。
完整的端到端行为验证可参考 tests/sandbox/test_session_sinks.py 中的相关测试,它们是理解各 Sink 契约的最直接示例。
七、选型建议与注意事项
按目标选 Sink:
| 需求 | 推荐 Sink |
|---|---|
| 进程内实时消费事件(告警、监控、测试断言) | CallbackSink |
| 宿主机侧留存 JSONL 审计日志 | JsonlOutboxSink |
| 日志随工作区走、跨 Docker/Modal 无需挂载卷 | WorkspaceJsonlSink |
| 转发到远程/本地代理服务 | HttpProxySink |
| 同时投递多个目标且保持顺序 | ChainedSink |
注意事项:
- 敏感数据:exec 输出默认不进入事件,需要时显式开启
include_exec_output并注意max_stdout_chars/max_stderr_chars的截断上限(默认各 8000 字符);write 事件永不包含文件字节。 - 时序窗口:会话启动早期与
shutdown收尾阶段,工作区可能不可写,WorkspaceJsonlSink会基于running()检查与阶段判断自动跳过刷新,避免误报。 - 绑定要求:
CallbackSink/WorkspaceJsonlSink等实现了SandboxSessionBoundSink的 Sink 必须经由SandboxSession(或带插桩的客户端)绑定后才生效,直接调用handle()会因未绑定而报错或空操作。 - 错误策略:生产环境建议
JsonlOutboxSink/HttpProxySink保持默认的best_effort+log,而把需要强保证的链路(如审计合规)配为sync+raise并辅以HttpProxySink的spool_path做失败溢出。
八、小结
Sinks 是 openai-agents-python 沙箱审计体系的外向出口:EventSink定义了统一的事件消费契约,CallbackSink、JsonlOutboxSink、WorkspaceJsonlSink、HttpProxySink分别覆盖进程内回调、宿主机文件、工作区文件与远程 HTTP 四类投递目标,ChainedSink提供顺序组合能力;配合DeliveryMode、OnErrorPolicy与EventPayloadPolicy,开发者可以按需组合出从实时监控到合规归档的完整事件管道。理解这套机制,是构建可观测、可审计、可追溯的沙箱 Agent 应用的关键一步。
【免费下载链接】openai-agents-pythonA lightweight, powerful framework for multi-agent workflows项目地址: https://gitcode.com/GitHub_Trending/op/openai-agents-python
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考