Skip to content

feat(file_edit): implement streaming file I/O with async operations, … - #15

Open
ykf173 wants to merge 1 commit into
wolf1069b:develop/agenticfrom
ykf173:feature/kaifeng.yan/fix_write_io_bug
Open

feat(file_edit): implement streaming file I/O with async operations, …#15
ykf173 wants to merge 1 commit into
wolf1069b:develop/agenticfrom
ykf173:feature/kaifeng.yan/fix_write_io_bug

Conversation

@ykf173

@ykf173 ykf173 commented Apr 13, 2026

Copy link
Copy Markdown
Collaborator

…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

…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

@gemini-code-assist gemini-code-assist Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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))

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

high

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.

Suggested change
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]

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

medium

Using only the first 8 characters of a UUID for a temporary file suffix increases the risk of collisions in high-concurrency environments. Consider using a longer portion or the full UUID to ensure uniqueness.

Suggested change
operation_id = str(uuid4())[:8]
operation_id = str(uuid4())

new_content = new_string
else:
# Use sublime_search for sophisticated replacement
from sublime_search import replace_content

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

medium

Scattered local imports can make the code harder to maintain and slightly slower if the function is called frequently. It is generally preferred to place all imports at the top of the module unless there is a specific reason for lazy loading (e.g., avoiding circular dependencies).

Comment on lines +359 to +400
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()

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

medium

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

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

medium

This local import is redundant as compute_unified_diff is already imported at the top of the file on line 25.

Leoyzen added a commit that referenced this pull request May 27, 2026
- 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
Leoyzen added a commit that referenced this pull request May 27, 2026
* 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
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant