ARTICLE DETAIL

资讯详情

深耕网站建设、视觉设计与SEO优化的一线实战洞察。

openai-agents-python 沙箱会话事件 Sinks 全解析:从回调、JSONL 落盘到 HTTP 代理的事件投递体系

openai-agents-python 沙箱会话事件 Sinks 全解析:从回调、JSONL 落盘到 HTTP 代理的事件投递体系 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()等底层原语工作区型 SinkWorkspaceJsonlSink正是通过这些原语把事件写进沙箱内部。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 需要声明四个属性属性类型含义namestr \| NoneSink 的可选名称便于日志与调试识别modeDeliveryMode投递模式见下文on_errorOnErrorPolicy出错时的处理策略payload_policyEventPayloadPolicy \| None事件负载裁剪策略控制敏感/大体积数据是否进入事件handle()是唯一的抽象方法负责把单个事件投递到目标。2. 关键类型别名模块顶部定义了两个受控的字面量类型sinks.pyDeliveryMode 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.pyruntime_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 将事件直接投递给用户提供的可调用对象同时支持同步与异步回调检测到协程对象时会awaitCallbackSink( 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), modesync)随后在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_keysTrue)并附加换行符保证每行独立可解析写入发生在线程池中asyncio.to_thread避免阻塞事件循环自动创建父目录mkdir(parentsTrue, exist_okTrue)在 POSIX 平台会尝试用fcntl.flock对文件加排他锁防止多进程并发追加时互相覆盖非 POSIX 平台如 Windows自动降级为不加锁默认best_effortlog的组合意味着日志写入失败不会拖垮沙箱主流程。测试 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_relpathlogs/events-{session_id}.jsonl工作区根目录下的相对路径支持轻量模板在bind()时展开ephemeralFalse为True时通过register_persist_workspace_skip_path()将日志路径排除出未来的工作区快照适合临时审计输出flush_every1每积累 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.HttpProxySinkPOST 事件到代理端点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按序组合多个 SinkChainedSink 用于把多个 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 的公共字段字段类型说明versionint事件格式版本默认1event_iduuid.UUID事件唯一 IDtsdatetimeUTC 时间戳session_iduuid.UUID所属会话 IDseqint会话内事件序号用于还原顺序opOpName操作名如exec、writephasestart \| finish阶段span_id/parent_span_id/trace_idstr \| None与 SDK 追踪联动的 span / trace 标识datadict操作相关的元数据路径、argv、耗时等SandboxSessionFinishEvent额外携带ok、duration_ms、error_code、error_type、error_message、error_retryable以及可选的stdout/stderr字符串原始字节输出stdout_bytes/stderr_bytes默认被excludeTrue排除在序列化之外仅用于 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(defaultFalse) # exec 输出默认关闭 max_stdout_chars: int Field(default8_000, ge0) max_stderr_chars: int Field(default8_000, ge0) include_write_len: bool Field(defaultTrue) # write 事件仅含尽力而为的字节数exec 输出默认关闭命令输出可能既嘈杂又敏感需要显式开启include_exec_outputTrue才会把截断后的stdout/stderr放进事件截断_safe_decode()以 UTF-8 容错解码errorsreplace按解码后字符串长度截断到max_stdout_chars/max_stderr_chars超出部分以…结尾保证 JSON 始终合法utils.pywrite 事件不含文件字节只尽力而为地统计写入长度通过 seekable 流的tell/seek估算见 utils.py绝不把文件内容放进审计事件。Instrumentation支持全局策略payload_policy与按操作细分的payload_policy_by_opdict[OpName, EventPayloadPolicy]后者可对特定操作覆盖全局配置manager.py。测试 test_sandbox_session_exec_emits_stdout_when_enabled 展示了开启include_exec_outputTrue后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_relpathPath(logs/events-{session_id}.jsonl), ephemeralTrue, ), # 同步回调事件失败实时可见 CallbackSink(on_event, modesync, on_errorraise), ) ], # 开启 exec 输出截断到默认 8000 字符 payload_policyEventPayloadPolicy(include_exec_outputTrue), ) # 以 UnixLocal 沙箱为例Docker 沙箱用法同理 inner UnixLocalSandboxSession.from_state( UnixLocalSandboxSessionState(manifestManifest(root/tmp/ws)) ) async with SandboxSession(inner, instrumentationinstrumentation) as session: await session.exec(echo hello from sandbox) await session.write(Path(note.txt), io.BytesIO(bdata))执行要点ChainedSink保证三个 Sink 按声明顺序依次收到每个事件CallbackSink使用modesync保证监控实时性JsonlOutboxSink/WorkspaceJsonlSink默认best_effort日志管道故障不会中断沙箱主流程ephemeralTrue让工作区内的审计日志不进入未来快照避免审计数据污染持久化工作区底层依赖 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_effortlog而把需要强保证的链路如审计合规配为syncraise并辅以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),仅供参考
返回列表