人工智能AI AgentAgent 框架大模型工具调用RAG提示工程强化学习【免费下载链接】agent-coreopenJiuwen agent-core可提供AI Agent开发、运行、调优与演进相关的全套SDK能力项目地址https://gitcode.com/openJiuwen/agent-core点击查看免费下载导读在 openJiuwen agent-core 的多智能体agent_teams架构中leader 通过Runner.run_agent_team_streaming(...)对外暴露流式接口但改造前这条流只能看到 leader 自己的 chunk——inprocess 模式下 spawn 出来的 teammate 走Runner.run_agent_team(memberTrue) → TeamAgent.invoke()路径其内部产出的思考、工具调用、文本输出全部被吞掉对上层CLI、监控、SDK 应用来说整个多成员协作就是一个黑盒。本文基于特性设计文档 F_02_member-attributed-streaming.md结合仓库源码与测试系统讲解成员归因流式输出的数据结构、数据流、核心决策与拒绝方案并给出消费端读取source_member/role的实战方式。读完本文你将掌握如何让 leader 的流同时携带所有 teammate 的 chunk、每个 chunk 如何标注归属成员与角色、为什么选择子类化 observer 钩子而不是消息总线或 envelope 包装以及该特性在 subprocess 模式的演进预留点。一、背景与要解决的问题1.1 改造前的数据黑洞Runner.run_agent_team_streaming(...)在改造前只能流出 leaderTeamAgent这一条 chunk 流。inprocess 模式下 spawn 出来的 teammate 走Runner.run_agent_team(memberTrue) → TeamAgent.invoke()路径invoke内部虽然创建了stream_queue并把 DeepAgent 产出的每个 chunk 都入了队但只把last_result当返回值——所有中间 chunk 在 teammate 自己的 invoke 主循环里被吞掉外层调用方完全看不到 teammate 的思考 / 工具调用 / 文本输出。1.2 对上层的影响对上层CLI、监控、调用 SDK 的应用来说多 teammate 协同就是一个黑盒——只能从 leader 的 chunk 看到协调动作看不到子任务现场。要做端到端可观测、做实时 UI 反馈就必须让 leader 的 streaming 同时携带 teammate 的 chunk且每个 chunk 标注归属成员。1.3 本期目标范围本期目标限定为 inprocess 模式同一进程、同一 event loop可以走对象引用subprocess 模式留扩展点。对应源码入口见 team_runner.py 的run_agent_team_streaming与内部 _run_team_member_streaming。二、数据结构TeamOutputSchema 子类2.1 新增的 TeamOutputSchema成员归因的核心数据结构是新增的TeamOutputSchema放在 agent_teams/schema/stream.pyclass TeamOutputSchema(OutputSchema): OutputSchema extended with the source-member identity and role. source_member: str | None None role: TeamRole | None None classmethod def from_output( cls, base: OutputSchema, *, source_member: str | None, role: TeamRole | None None, ) - TeamOutputSchema: Build a tagged team chunk from a plain OutputSchema instance. return cls( typebase.type, indexbase.index, payloadbase.payload, source_membersource_member, rolerole, )它继承自 core 层的 OutputSchematype/index/payload三个字段的 pydanticBaseModel新增两个字段source_member: str | None产出该 chunk 的成员member_namerole: TeamRole | None该成员在团队中的角色。from_output类方法接收一个普通OutputSchema实例并返回一个新的TeamOutputSchema——不修改原对象保证 DeepAgent 内部持有的对象引用不受影响测试 test_tag_chunk_upgrades_plain_outputschema 明确断言原始 chunk 不能被 mutate。2.2 为什么子类化而不是改 core 层设计上不动 core 层OutputSchema——core 跨子系统共享不应被 team-specific 字段污染。子类对isinstance(x, OutputSchema)透明现有消费端访问chunk.type / chunk.payload不变。TeamOutputSchema已通过 schema/init.py 和 agent_teams/init.py 导出为公开 API。2.3 role 字段的语义role字段使用 TeamRole 枚举核心值包括LEADER/TEAMMATE/HUMAN_AGENT以及PASSIVE_HUMAN、BRIDGE_AGENT、WORKER、EXTERNAL_CLI等扩展角色。TeamRole是str枚举pydantic 序列化得到字符串值跨进程兼容。消费端拿到 chunk 后能立即区分协调指令 vs 子任务输出 vs 人类输入的语义无需维护member_name → role的外部映射。2.4 团队事件标记识别同文件中还提供了辅助函数 is_team_event_marker当 chunk 的payload.event_type以team.开头如team.idle/team.completed/team.interact.failed时返回True。团队标记与模型输出走同一条流保证流式消费者按序读取非流式调用方若把流规约为最后一个产出可用它跳过框架事件。三、数据流teammate chunk 如何进入 leader 的流3.1 整体数据流图inprocess 模式Teammate DeepAgent │ run_streaming() yields raw OutputSchema chunk ▼ Teammate StreamController._forward_outputs │ _tag_chunk(): 升级为 TeamOutputSchema(source_member, role) ├──────► teammate.stream_queue.put(chunk) [本地] └──────► for observer in _chunk_observers: [fan-out] forward(chunk) │ ▼ Leader StreamController.stream_queue.put(chunk) │ ▼ TeamAgent.stream() yield chunk │ ▼ Runner.run_agent_team_streaming yields chunk当前StreamController中的实际循环入口是 _forward_outputs负责从 runtime 的outputs()拉流、_tag_chunk打标、入本地队列并 fan-out 给所有 observer。leader 自己的 chunk 走相同的_tag_chunk路径observer 列表为空直接进自己 queue。所有出口 chunk 都是TeamOutputSchema没有leader 走哪条 / teammate 走哪条的代码分支。3.2 打标逻辑_tag_chunk_tag_chunk 的规则是普通OutputSchema→ 用TeamOutputSchema.from_output(...)升级携带当前成员的member_name与role已经是TeamOutputSchema且source_member/role匹配当前成员 → 原样返回避免无谓拷贝已经是TeamOutputSchema但成员不匹配 →model_copy(update...)重新打标原对象不动非OutputSchema的 chunk → 原样透传。对应测试见 test_stream_controller.py升级、透传、重写、非 OutputSchema 四种分支全覆盖。3.3 成员身份的注入成员的member_name与role通过blueprint_getter回调获取self._get_blueprint()→bp.member_name/bp.role对应 StreamController 构造参数。blueprint_getter返回 TeamAgentBlueprint拿不到 blueprint 时_member_name()返回None打标自然回退为source_memberNone。四、核心决策与底层实现决策 1派生子类不在 core 层加字段把source_member字段放在TeamOutputSchema(OutputSchema)而不是OutputSchema上。子类化让消费端透明同时保持 core/session/stream/base.py 的纯净——source_member只对 team 子系统有意义single_agent / harness 路径继续 yield 原生OutputSchema。决策 2observer hook 而不是消息总线leader 与 teammate 共享同一个 event loopinprocess 同进程直接走对象引用最简单。在 StreamController 上加了一对方法def add_chunk_observer(self, cb: ChunkObserver) - None: Register a chunk observer fired after each chunk is tagged. self._chunk_observers.append(cb) def remove_chunk_observer(self, cb: ChunkObserver) - None: Detach a previously-registered observer; idempotent. with contextlib.suppress(ValueError): self._chunk_observers.remove(cb)ChunkObserver Callable[[OutputSchema], Awaitable[None]]类型别名。每次 chunk 入队时同步 fan-out 给所有 observerobserver 抛异常自动 detach不阻塞主流由team_logger.exception记录——见 _forward_outputs。observer 反注册幂等有测试覆盖test_remove_chunk_observer_is_idempotent。决策 3SpawnManager 是单点 wiring 责任方teammate 的StreamController不知道 leader 是谁leader 的StreamController不知道 teammate 在哪。SpawnManager._wire_inprocess_chunk_forward(handle) 是唯一持有两边引用的地方def _wire_inprocess_chunk_forward(self, handle: InProcessSpawnHandle) - None: Forward an in-process teammates stream chunks into leaders queue. leader self._get_team_agent() if leader is None or handle.agent_ref is None: return leader_sc leader.stream_controller teammate_sc handle.agent_ref.stream_controller async def _forward(chunk: OutputSchema) - None: # Drop silently if the leaders queue is not set up yet or # has already been torn down — buffering would invert the #>return TeamOutputSchema( typemessage, index0, payload{ event_type: team.runtime_ready, team_name: team_name, session_id: session_id, activation_kind: action_kind.value, }, source_memberleader_member_name, roleleader_role, )leadermember_name通过activation.agent.blueprint.member_name取见 run_agent_team_streaming。注意这里采用惰性导入from openjiuwen.agent_teams.schema.stream import TeamOutputSchema保证 agent_teams 不进子进程的 bootstrap 路径。五、拒绝的方案设计权衡记录方案 A在 core 层OutputSchema上加source_member字段拒绝理由OutputSchema是跨 single_agent / harness / agent_teams 共享的数据结构。source_member只对 team 层有意义给 core 加 team-specific 字段是概念污染违反分层边界。子类化在isinstance关系上等价、消费端无感知是更干净的扩展。方案 B包 envelopeTeamStreamChunk(member_name, role, chunk)拒绝理由envelope 包装让所有现有消费端必须改chunk.payload → env.chunk.payload是大面积破坏性变更继承可以做到同样的语义而对老代码透明仍是OutputSchema子类。方案 Cmessager 总线广播 chunk拒绝理由inprocess 同进程同 event loop对象引用一步到位没有理由绕一圈进 pub/sub。messager 适合跨进程subprocess 模式本期不需要。如果未来要支持 subprocess扩展点已留好——add_chunk_observer的接口和TeamOutputSchema结构不变subprocess 端把 chunk publish 到TeamTopic.STREAM_CHUNK、leader 端订阅并反序列化 put 到自己 queue 即可。方案 Dteammate 主路径改成走 streaming不再丢 chunk拒绝理由Runner.run_agent_team(memberTrue)现在的 invoke 路径与 spawn 工具调用约定深度耦合改它会牵连子进程入口与from_spawn_payloadwire 协议。observer fan-out 实现同样的语义但只新增数据通路、不破坏既有契约。六、验证测试基线设计文档记录的测试基线在仓库中可逐条对应tests/unit_tests/agent_teams/test_stream_controller.py文档记录16 passed其中_tag_chunk四分支升级 / 同身份透传 / 异身份重写 / 非 OutputSchema 透传、_forward_outputs打标 fan-outtest_forward_outputs_tags_and_fans_out_to_observers、observer 异常自动 detachtest_observer_exception_auto_detaches_and_does_not_block_stream、observer 反注册幂等、teammate → leader queue 端到端数据流test_teammate_chunks_reach_leader_queue_via_forward_observer断言两个 chunk 均带source_member teammate_m且role TeamRole.TEAMMATE、leader queue 为 None 时丢弃。tests/unit_tests/agent_teams/test_spawn_manager_chunk_forward.py文档记录3 passed覆盖_wire_inprocess_chunk_forward注入路径test_wire_forward_routes_teammate_chunk_to_leader_queue、cleanup_teammate反注册test_cleanup_detaches_forward_observer、leader / agent_ref 缺失时 wire no-optest_wire_skips_when_leader_or_agent_ref_missing。现有 streaming 用例无需修改向后兼容。七、消费端实战如何读取成员归因的 chunk7.1 标准消费循环Runner.run_agent_team_streaming(...)产出的每个 chunk 都是TeamOutputSchema消费端可以这样读取from openjiuwen.agent_teams.schema.stream import TeamOutputSchema, is_team_event_marker from openjiuwen.core.runner.runner import Runner async for chunk in Runner.run_agent_team_streaming(team_spec, inputs): assert isinstance(chunk, TeamOutputSchema) # team 路径下全部打标 member chunk.source_member # 产出该 chunk 的成员名 role chunk.role # LEADER / TEAMMATE / ... if is_team_event_marker(chunk): # team.runtime_ready / team.idle / team.completed 等框架事件 continue # member / role 分支区分协调指令 vs 子任务输出 vs 人类输入 ...注意几个事实约束旧消费端chunk.type / chunk.payload访问路径不变isinstance透明只有 leader 路径run_agent_team_streaming非 member 分支上能看到汇总后的全团队流memberTrue分支_run_team_member_streaming是 teammate 自身视角的流团队完成标记team.completed严格排在Nonesentinel 之前emit_completion_and_close所以流式消费者会先读到完成信号再看到流结束。7.2 团队成员生命周期事件除了 chunk 归因同一控制器还会在流上产出两种团队事件标记同样是TeamOutputSchemapayload.event_type以team.开头team.idle团队进入静默防抖窗口_TEAM_IDLE_DEBOUNCE_SECONDS 2.0秒内所有成员保持静止、且任务板无遗留任务时触发payload 携带整队 roster 快照member_countmembers见 emit_team_idleteam.completed团队完成携带member_count/task_count。八、已知遗留与演进方向文档明确列出了三期后续工作均已在源码中留下对应位点subprocess 模式的 chunk 转发本期不实现。扩展方向teammate 进程 publish chunkTeamOutputSchema.model_dump()到TeamTopic.STREAM_CHUNK、leader 进程订阅并反序列化为TeamOutputSchema后 put 到自己 queue。add_chunk_observer与TeamOutputSchema数据结构无需改动。要点跨进程序列化 chunk 是性能开销点可加TeamAgentSpec.stream_member_chunks: bool让用户按需开关。CLI / 示例的来源展示cli/stream_renderer.py当前未利用source_member做着色或前缀后续可加一个 per-member 颜色映射把不同成员的 chunk 在 TUI 中可视化区分。这是渲染优化不是协议变更。相关 TUI 设计可见 S_15_cli-tui.md。chunk_observer 的有界化当前 fan-out 是for ob in list(...): await ob(chunk)串行调用。生产环境 observer 只有 forward 一个无阻塞如果未来挂多个高延迟 observer要考虑改成并发 gather 单 observer 超时熔断。九、涉及文件速查作用文件路径数据结构TeamOutputSchema/is_team_event_markeropenjiuwen/agent_teams/schema/stream.py角色枚举TeamRoleopenjiuwen/agent_teams/schema/team.py打标 observer 数据流主循环openjiuwen/agent_teams/agent/stream_controller.pywiring 责任方_wire_inprocess_chunk_forward/cleanup_teammateopenjiuwen/agent_teams/agent/spawn_manager.pyinprocess 句柄chunk_forward字段openjiuwen/agent_teams/spawn/inprocess_handle.pyready chunk 升级 流式入口openjiuwen/core/runner/team_runner.pycore 基类OutputSchemaopenjiuwen/core/session/stream/base.py控制器测试tests/unit_tests/agent_teams/test_stream_controller.pywiring / cleanup 测试tests/unit_tests/agent_teams/test_spawn_manager_chunk_forward.py十、小结Member-Attributed Streaming 用一个TeamOutputSchema子类 一个add_chunk_observer钩子 SpawnManager 单点 wiring三件套以最小破坏面解决了多智能体流式输出的可观测性痛点所有出口 chunk 统一带source_member与role旧消费端零改动inprocess 模式下通过对象引用单向转发、无消息总线开销、无反向耦合同时为 subprocess 模式预留了TeamTopic.STREAM_CHUNK演进路径。这套设计的取舍子类化 vs envelope、observer vs 总线、惰性队列 vs 缓冲对任何需要在共享流协议之上叠加生产者身份的场景都有直接参考价值。赞分享人工智能AI AgentAgent 框架大模型工具调用RAG提示工程强化学习【免费下载链接】agent-coreopenJiuwen agent-core可提供AI Agent开发、运行、调优与演进相关的全套SDK能力项目地址https://gitcode.com/openJiuwen/agent-core点击查看免费下载相关推荐POCO多线程性能调优团队协作角色与流程POCO多线程性能调优团队协作角色与流程 引言 在当今软件开发中多线程技术是提升应用性能的关键手段之一。然而多线程编程也带来了诸多挑战如线程安全、死锁、后端网络/通信数据库密码学Web框架AWS S3安全配置DevOps Interview Guide中的bucket保护策略AWS S3安全配置DevOps Interview Guide中的bucket保护策略 在当今云原生时代AWS S3作为核心存储服务其安全配置直接关系到文档知识库DevOpsopenJiuwen agent-core HITT 模式实战让真实人类以团队成员身份参与 Agent 团队协作openJiuwen agent core HITT 模式实战让真实人类以团队成员身份参与 Agent 团队协作 Human In The TeamHITT人工智能AI AgentAgent 框架大模型工具调用RAG提示工程强化学习上一篇从 PCAP 到 IOC基于 Anthropic-Cybersecurity-Skills 的网络流量事件分析工具 API 参考实战指南下一篇Tabby 自托管 AI 编码助手Docker 一键部署、CLI 参数解析与源码构建完整指南创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考