feat(file_edit): implement streaming file I/O with async operations, … - #15
feat(file_edit): implement streaming file I/O with async operations, …#15ykf173 wants to merge 1 commit into
Conversation
…atomic writes, and progress reporting Implements streaming file editing capabilities based on write tool optimization proposal: Features: - Async non-blocking I/O using aiofiles (prevents UI freezing) - Atomic writes via temp file + rename pattern (prevents data corruption) - Real-time progress events via ToolCallProgressEvent - Configurable timeout protection (asyncio.timeout) - Chunked streaming for large files (default 64KB chunks) APIs: - StreamingFileEditor: Main class with async I/O and progress - StreamingWriteTool: Sync-style wrapper for simple usage - streaming_write_file()/streaming_edit_file(): Convenience functions Testing: - 24 comprehensive tests covering basic ops, progress events, atomic writes, timeouts, errors, edge cases - All tests passing Documentation: - Full docstrings and type hints - Implementation report (STREAMING_WRITE_IMPLEMENTATION_REPORT.md) - Interactive demo script (examples/streaming_file_write_demo.py) Backward compatible: Legacy edit_file_tool unchanged
There was a problem hiding this comment.
Code Review
This pull request implements a streaming file editing system featuring async I/O via aiofiles, atomic writes, and real-time progress events. The reviewer provided several actionable suggestions to improve the robustness and idiomatic quality of the code, such as using os.replace for cross-platform atomic renames, utilizing full UUIDs for temporary files to prevent collisions, and employing async with contexts for safer file handling. Additionally, feedback was given regarding the consolidation of local imports and the removal of redundant ones to enhance maintainability.
| # Atomic rename (run in thread to avoid blocking) | ||
| import os | ||
|
|
||
| await asyncio.to_thread(os.rename, str(temp_path), str(path)) |
There was a problem hiding this comment.
os.rename can fail on Windows if the destination file already exists, raising a FileExistsError. os.replace is the cross-platform standard for atomic replacement that works even if the target exists.
| await asyncio.to_thread(os.rename, str(temp_path), str(path)) | |
| await asyncio.to_thread(os.replace, str(temp_path), str(path)) |
| Yields: | ||
| ToolCallProgressEvent with progress updates | ||
| """ | ||
| operation_id = str(uuid4())[:8] |
There was a problem hiding this comment.
| new_content = new_string | ||
| else: | ||
| # Use sublime_search for sophisticated replacement | ||
| from sublime_search import replace_content |
There was a problem hiding this comment.
| f = await aiofiles.open(temp_path, mode="wb") | ||
| try: | ||
| for chunk_num in range(total_chunks): | ||
| start_idx = chunk_num * self.chunk_size | ||
| end_idx = min(start_idx + self.chunk_size, total_size) | ||
| chunk = content_bytes[start_idx:end_idx] | ||
|
|
||
| await f.write(chunk) | ||
| bytes_written += len(chunk) | ||
|
|
||
| # Calculate progress metrics | ||
| elapsed = time.monotonic() - start_time | ||
| bytes_per_second = bytes_written / elapsed if elapsed > 0 else 0 | ||
| percentage = bytes_written / total_size if total_size > 0 else 1.0 | ||
| remaining_bytes = total_size - bytes_written | ||
| estimated_seconds_remaining = ( | ||
| remaining_bytes / bytes_per_second if bytes_per_second > 0 else 0 | ||
| ) | ||
|
|
||
| # Yield progress event every few chunks or at start/end | ||
| if chunk_num % 4 == 0 or chunk_num == total_chunks - 1: | ||
| progress_items: list[ToolCallContentItem] = [ | ||
| LocationContentItem(path=str(path), line=0), | ||
| TextContentItem( | ||
| text=( | ||
| f"Writing: {bytes_written}/{total_size} bytes " | ||
| f"({percentage * 100:.1f}%) " | ||
| f"@ {bytes_per_second / 1024:.1f} KB/s" | ||
| ) | ||
| ), | ||
| ] | ||
|
|
||
| yield ToolCallProgressEvent( | ||
| tool_call_id=tool_call_id, | ||
| status="in_progress", | ||
| title=f"Writing {path.name}: {percentage * 100:.0f}%", | ||
| items=progress_items, | ||
| progress=chunk_num + 1, | ||
| total=total_chunks, | ||
| ) | ||
| finally: | ||
| await f.close() |
There was a problem hiding this comment.
Using async with for aiofiles.open is more idiomatic and safer than manually calling close() in a finally block. It ensures the file is closed correctly even if an exception occurs during the iteration.
try:
async with aiofiles.open(temp_path, mode="wb") as f:
for chunk_num in range(total_chunks):
start_idx = chunk_num * self.chunk_size
end_idx = min(start_idx + self.chunk_size, total_size)
chunk = content_bytes[start_idx:end_idx]
await f.write(chunk)
bytes_written += len(chunk)
# Calculate progress metrics
elapsed = time.monotonic() - start_time
bytes_per_second = bytes_written / elapsed if elapsed > 0 else 0
percentage = bytes_written / total_size if total_size > 0 else 1.0
remaining_bytes = total_size - bytes_written
estimated_seconds_remaining = (
remaining_bytes / bytes_per_second if bytes_per_second > 0 else 0
)
# Yield progress event every few chunks or at start/end
if chunk_num % 4 == 0 or chunk_num == total_chunks - 1:
progress_items: list[ToolCallContentItem] = [
LocationContentItem(path=str(path), line=0),
TextContentItem(
text=(
f"Writing: {bytes_written}/{total_size} bytes "
f"({percentage * 100:.1f}%) "
f"@ {bytes_per_second / 1024:.1f} KB/s"
)
),
]
yield ToolCallProgressEvent(
tool_call_id=tool_call_id,
status="in_progress",
title=f"Writing {path.name}: {percentage * 100:.0f}%",
items=progress_items,
progress=chunk_num + 1,
total=total_chunks,
)| for item in event.items: | ||
| if hasattr(item, "old_text") and hasattr(item, "new_text"): | ||
| # This is a DiffContentItem | ||
| from agentpool.utils.diffs import compute_unified_diff |
* docs(rfc): add RFC-0034 ACP Session Config Options 统一化 新增 RFC-0034,提案升级 AgentPool ACP Server 的 Session Config Options 透出逻辑,使 Zed 等 ACP 兼容 IDE 能够选择模型和切换 Agent Role。 主要内容: - 识别 4 个 GAP:Agent Role 未透出(P0)、ACP/OpenCode model list 数据来源不一致(P1)、/mode 路由硬编码(P1)、get_session_mode_state 过滤过严(P2) - 分析 3 个方案,推荐选项 2(三阶段统一化) - 技术设计:build_model_state_for_acp()、get_agent_role_config_option()、 _swap_session_agent() 及 OpenCode /mode 路由动态修复 🤖 Generated with [Qoder][https://qoder.com] * docs(rfc): 根据 review 反馈修正 RFC-0034 - 修正 model fallback 逻辑:strict fallback(configured 存在时只用 configured) - 修正 agent_role current_value:使用 agent.name 而非 pool.main_agent.name - 重写 _swap_session_agent():委托 session.switch_active_agent() + _session_agent_locks 保护 - 增加 session._task_lock 协调:拒绝 active prompt 期间的 swap - 增加 pool.manifest null check 保护 - 修正 list_modes() null safety:state.agent 为 None 时返回默认值 - 明确开放问题 Q2/Q3 的决策:对话历史不继承、current_value 已修复 - 更新决策记录:增加 session mutation 复用、锁保护、task_lock 协调、对话历史决策 - 更新 Phase 2 实施计划:增加 Zed 预验证、并发测试、current_value 测试 - 调整工作量估算:~260 行 → ~240 行 * docs(rfc): RFC-0034 新增 Phase 0 — ACP Configurable LLM Providers 适配 ACP PR #648 (Configurable LLM Providers) 已 MERGED,引入 providers/list、providers/set、providers/disable 三个方法族, 允许客户端发现和覆盖 agent 的 LLM 请求路由。 主要更新: - 新增 GAP 5 (P0): providers/* 完全未实现 - 新增目标 G7: 实现 ACP providers/* 协议方法 - 新增 Phase 0: ProviderRouter 实现 + schema 类型定义 + ACP 请求处理器 + AgentCapabilities.providers 声明 - 修订 Phase 1: build_model_state_for_acp() 接受 provider_router 参数,过滤被禁用 provider 下的模型 - 更新架构概览图: 传输层(providers)与应用层(session config)分层 - 更新里程碑: 四阶段实施,Phase 0 优先于 Phase 1 - 新增开放问题 6/7/8: providers 对已运行 session 的影响、 provider 路由覆盖与 agent 初始化兼容、SessionModelState 中是否携带 provider 关联信息 - 新增决策记录: providers/set 保守策略、从 model_variants 派生 ProviderInfo、provider_router 参数解耦 🤖 Generated with [Qoder][https://qoder.com] * docs(rfc): 优化 RFC-0034 — 补充 Zed 源码级兼容性分析 基于 Zed 源码调研(crates/agent_ui/src/config_options.rs、profile_selector.rs、 agent_servers/src/acp.rs)的关键发现: 1. Zed 渲染所有 config_options 为独立 UI 按钮,agent_role 可正确显示和点击 2. first_config_option_id() 仅返回同 category 的第一个 option,键盘快捷键 可能冲突 — 标记为已知限制(NG7) 3. Zed ProfileSelector 完全独立于 ACP,使用本地 AgentSettings.profiles 4. Zed 当前完全不支持 providers/* 协议(Phase 0 暂无 Zed UI 入口) RFC 更新内容: - 新增 Zed IDE 渲染行为小节(源码级证据) - 新增 Zed 兼容性分析总结表 - 更新非目标 NG7:键盘快捷键冲突为已知限制 - 更新开放问题 5/6/7/8/9,标记 Zed 调研结论 - 更新 Phase 2 预验证:明确键盘限制和排序建议 - 更新向后兼容保证表:添加 category 冲突行 - 更新决策记录:补充 Zed 调研证据 🤖 Generated with [Qoder][https://qoder.com] * feat(acp): implement RFC-0034 ACP Session Config Options unification Phase 0: ACP Configurable LLM Providers - Add providers/* protocol methods (providers/list, providers/set, providers/disable) - Add ProviderRouter with override/disable/capability tracking - Add providers field to AgentCapabilities and InitializeResponse Phase 1: Shared Model List Logic - Add build_model_state_for_acp() with configured-first, tokonomics-fallback - Invert get_session_model_state() to use configured variants first Phase 2: Agent Role Config Option - Add get_agent_role_config_option() exposing pool.all_agents - Add _swap_session_agent() with lock protection - Extend set_session_config_option() with agent_role handling Phase 3: OpenCode /mode Route Fix - Dynamic /mode route using agent.get_modes() Also includes RFC-0033 MCP over ACP support: - Add AcpMcpServer type and acp field to McpCapabilities - Add acp_mcp_servers parameter to AgentCapabilities.create() Tests: - 35 new tests across provider_router, model_state, agent_role, config_routes, and cross-protocol integration - Snapshot tests re-baselined * chore: remove RFC-0033 code from RFC-0034 branch Remove accidentally included RFC-0033 MCP-over-ACP implementation: - Delete acp_mcp_manager.py, acp_mcp_transport.py - Delete RFC-0033 tests (test_mcp.py, test_acp_mcp_*, test_mcp_integration) - Remove AcpMcpServer from mcp.py - Remove acp field from McpCapabilities - Remove acp_mcp_servers parameter from AgentCapabilities.create() - Remove acp_mcp_servers parameter from InitializeResponse.create() - Remove RFC-0033 handler code from acp_agent.py Keep RFC-0034 changes intact: - providers/* protocol methods - ProviderRouter with override/disable - build_model_state_for_acp() configured-first logic - agent_role config option and swap - Dynamic /mode route * fix: address PR review comments for RFC-0034 - _swap_session_agent: Update session agent registry after swap (review #5) - get_agent_role_config_option: Use display_name with type-safe fallback, add description (review #7) - list_modes: Use mode.id for programmatic identifiers, add explicit None guard (review #4, #8) - Update tests to match new behavior * fix(agent): use model_variants in get_modes() instead of tokonomics Agent.get_modes() was calling get_available_models() which returns all tokonomics-discovered models (2000+). Now it checks configured model_variants first and only falls back to tokonomics when no variants are configured. Fixes the issue where config_options model selector showed thousands of models instead of the configured variants. * fix(agent): track model variant name to fix Zed Unknown display When using model_variants, get_modes() returned variant names as option ids but current_mode_id was the raw model identifier (e.g. openai:svc/glm-4.7). This caused Zed to display 'Unknown' because current_mode_id didn't match any available mode id. Fix: Add _current_model_variant field to Agent. When _set_mode() is called with a variant name, store it. get_modes() now uses _current_model_variant as current_mode_id so it matches the option ids. * fix(agent): set _current_model_variant on init when model is variant name Agent.__init__ resolves model string via _resolve_model_string(), but was not setting _current_model_variant. This caused get_modes() to fall back to self.model_name (the raw model identifier) on initial load, showing 'Unknown' in Zed until _set_mode() was called. Fix: Also track variant name in __init__ when model string matches a model_variants key. * fix(agent): use actual model identifier as mode id, variant name as display name Redesign model config option to use actual model identifiers: - id/value: actual model identifier (e.g. openai:svc/glm-4.7) - name: variant name (e.g. glm47) for display - current_mode_id: actual model identifier This ensures currentValue matches option values in Zed's config option selector, fixing the 'Unknown' display issue. _set_mode() now supports both actual model identifiers and variant names by reverse-lookup from manifest model_variants. * fix(agent): align get_modes() id format with model_name Use config.get_model().system:model_name for option ids instead of config.identifier, ensuring currentValue matches option values. Root cause: model_name returns pydantic-ai system:model_name format (e.g., 'openai:svc/glm-4.7') while config.identifier returns full provider format (e.g., 'openai-chat:svc/glm-4.7'), causing mismatch in Zed's model selector dropdown. * fix: address PR #37 review comments (round 2) - providers/set & providers/disable: 兼容 id 字段(Comment #12, #13) - provider_router: 防御性初始化 + 未知 provider 静默禁用(Comment #14, #15) - model_utils: 先过滤 raw toko_models(更准确),current_model 不在列表时插入(Comment #16) - .gitignore: 添加 .omo/(Comment #18)
…atomic writes, and progress reporting
Implements streaming file editing capabilities based on write tool optimization proposal:
Features:
APIs:
Testing:
Documentation:
Backward compatible: Legacy edit_file_tool unchanged