# ========= Copyright 2025-2026 @ Eigent.ai All Rights Reserved. ========= # Licensed under the Apache License, Version 2.0 (the "License"); # you may not use this file except in compliance with the License. # You may obtain a copy of the License at # # http://www.apache.org/licenses/LICENSE-2.0 # # Unless required by applicable law or agreed to in writing, software # distributed under the License is distributed on an "AS IS" BASIS, # WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. # See the License for the specific language governing permissions and # limitations under the License. # ========= Copyright 2025-2026 @ Eigent.ai All Rights Reserved. ========= import os import uuid from dataclasses import dataclass from pathlib import Path from typing import Literal from camel.toolkits import FileToolkit as BaseFileToolkit from app.agent.toolkit.abstract_toolkit import AbstractToolkit from app.component.environment import env from app.run_context import RunContext from app.run_runtime.tool_checkpoint import ( ToolInvocationNotDispatchedError, get_current_tool_checkpoint, ) from app.service.task import ( ActionWriteFileData, Agents, get_task_lock, process_task, ) from app.utils.listen.toolkit_listen import ( _safe_put_queue, auto_listen_toolkit, listen_toolkit, ) from app.utils.space_overlay_client import ( path_write_lock, post_overlay_write, relative_to_artifact_root, relative_to_workdir, run_context_for_task, sha256_of_file, should_record_overlay, ) from app.workspace_git import get_default_workspace_mutation_service @dataclass(frozen=True) class OverlayWriteContext: run_context: RunContext rel_path: str target_path: Path @dataclass(frozen=True) class PendingOverlayWrite: run_context: RunContext rel_path: str target_path: Path base_hash: str | None status: Literal["added", "modified"] file_hash: str | None size: int mode: int @auto_listen_toolkit(BaseFileToolkit) class FileToolkit(BaseFileToolkit, AbstractToolkit): agent_name: str = Agents.document_agent def __init__( self, api_task_id: str, working_directory: str | None = None, timeout: float | None = None, default_encoding: str = "utf-8", backup_enabled: bool = True, ) -> None: if working_directory is None: working_directory = env( "file_save_path", os.path.expanduser("~/Downloads") ) super().__init__( working_directory, timeout, default_encoding, backup_enabled ) self.api_task_id = api_task_id def _prepare_file_write( self, mutation_service, *, run_context: RunContext, filename: str, operation_request_id: str, actor_id: str, trigger: str, ): """Prepare Git state before the file operation can touch its target.""" try: return mutation_service.prepare_file_write( context=run_context, filename=filename, operation_request_id=operation_request_id, actor_id=actor_id, trigger=trigger, ) except Exception as error: # Workspace snapshot/admission happens before CAMEL invokes the # underlying file operation. Preserve that boundary so an overlay # conflict is a known, model-visible failure rather than an # ambiguous external side effect that poisons the whole Run. error_code = getattr(error, "code", None) detail = f"{error_code}: {error}" if error_code else str(error) raise ToolInvocationNotDispatchedError( "File mutation was not started because its Workspace could " f"not be prepared: {detail}" ) from error def _overlay_write_context( self, filename: str ) -> OverlayWriteContext | None: context = run_context_for_task(self.api_task_id) if context is None: return None resolved = relative_to_workdir(context, filename) if resolved is None: return None rel_path, target_path = resolved if not should_record_overlay(context, target_path): return None return OverlayWriteContext(context, rel_path, target_path) @listen_toolkit( BaseFileToolkit.write_to_file, lambda _, title, content, filename, encoding=None, use_latex=False: ( f"write content to file: {filename} with encoding: {encoding} and use_latex: {use_latex}" ), ) def write_to_file( self, title: str, content: str | list[list[str]], filename: str, encoding: str | None = None, use_latex: bool = False, ) -> str: run_context = run_context_for_task(self.api_task_id) prepared_write = None relative_path = None operation_request_id = None mutation_service = None if run_context is not None: checkpoint = get_current_tool_checkpoint() operation_request_id = ( checkpoint.tool_call_id if checkpoint is not None else f"local-file-write:{uuid.uuid4().hex}" ) mutation_service = get_default_workspace_mutation_service() prepared_write = self._prepare_file_write( mutation_service, run_context=run_context, filename=filename, operation_request_id=operation_request_id, actor_id=self.agent_name, trigger="filesystem.write", ) if prepared_write is not None: res = super().write_to_file( title, content, str(prepared_write.target_path), encoding, use_latex, ) if "Content successfully written to file: " in res: assert operation_request_id is not None assert mutation_service is not None mutation_service.complete_file_write( prepared_write, operation_request_id=operation_request_id, actor_id=self.agent_name, trigger="filesystem.write", ) res = ( "Content successfully written to file: " f"{prepared_write.relative_path}" ) relative_path = prepared_write.relative_path else: res = self._write_legacy_or_cloud_overlay( title=title, content=content, filename=filename, encoding=encoding, use_latex=use_latex, ) if "Content successfully written to file: " in res: written_path = res.replace( "Content successfully written to file: ", "" ) if relative_path is None and run_context is not None: relative_path = relative_to_artifact_root( run_context, written_path ) task_lock = get_task_lock(self.api_task_id) # Capture ContextVar value before creating async task current_process_task_id = process_task.get("") # Use _safe_put_queue to handle both sync and async contexts _safe_put_queue( task_lock, ActionWriteFileData( process_task_id=current_process_task_id, data=written_path, relative_path=relative_path, ), ) return res def edit_file( self, file_path: str, old_content: str, new_content: str, ) -> str: """Edit one file inside the Run-owned Agent worktree. CAMEL's base implementation resolves relative paths against the toolkit working directory. That directory is the visible Space for compatibility sessions, so calling it directly would let follow-up edits bypass the Run commit. Route the edit through the same durable exact-path mutation used by ``write_to_file`` instead. """ run_context = run_context_for_task(self.api_task_id) prepared_write = None relative_path = None operation_request_id = None mutation_service = None if run_context is not None: checkpoint = get_current_tool_checkpoint() operation_request_id = ( checkpoint.tool_call_id if checkpoint is not None else f"local-file-edit:{uuid.uuid4().hex}" ) mutation_service = get_default_workspace_mutation_service() prepared_write = self._prepare_file_write( mutation_service, run_context=run_context, filename=file_path, operation_request_id=operation_request_id, actor_id=self.agent_name, trigger="filesystem.edit", ) if prepared_write is not None: res = super().edit_file( str(prepared_write.target_path), old_content, new_content, ) if res.startswith("Successfully edited "): assert operation_request_id is not None assert mutation_service is not None mutation_service.complete_file_write( prepared_write, operation_request_id=operation_request_id, actor_id=self.agent_name, trigger="filesystem.edit", ) relative_path = prepared_write.relative_path res = f"Successfully edited {relative_path}" else: res = super().edit_file(file_path, old_content, new_content) if res.startswith("Successfully edited "): edited_path = res.removeprefix("Successfully edited ") if relative_path is None and run_context is not None: relative_path = relative_to_artifact_root( run_context, edited_path ) task_lock = get_task_lock(self.api_task_id) current_process_task_id = process_task.get("") _safe_put_queue( task_lock, ActionWriteFileData( process_task_id=current_process_task_id, data=edited_path, relative_path=relative_path, ), ) return res def _write_legacy_or_cloud_overlay( self, *, title: str, content: str | list[list[str]], filename: str, encoding: str | None, use_latex: bool, ) -> str: overlay_context = self._overlay_write_context(filename) if overlay_context is None: res = super().write_to_file( title, content, filename, encoding, use_latex ) else: pending_overlay_write: PendingOverlayWrite | None = None with path_write_lock( overlay_context.run_context.space_id, overlay_context.run_context.project_id, overlay_context.run_context.run_id, overlay_context.rel_path, ): existed = overlay_context.target_path.exists() base_hash = sha256_of_file(overlay_context.target_path) res = super().write_to_file( title, content, filename, encoding, use_latex ) if "Content successfully written to file: " in res: written_path = Path( res.replace( "Content successfully written to file: ", "" ) ) if not written_path.is_absolute(): written_path = ( Path(self.working_directory) / written_path ) written_hash = sha256_of_file(written_path) written_stat = written_path.stat() pending_overlay_write = PendingOverlayWrite( run_context=overlay_context.run_context, rel_path=overlay_context.rel_path, target_path=written_path, base_hash=base_hash, status="modified" if existed else "added", file_hash=written_hash, size=written_stat.st_size, mode=written_stat.st_mode, ) if pending_overlay_write is not None: post_overlay_write( pending_overlay_write.run_context, pending_overlay_write.rel_path, pending_overlay_write.target_path, base_hash=pending_overlay_write.base_hash, status=pending_overlay_write.status, file_hash=pending_overlay_write.file_hash, size=pending_overlay_write.size, mode=pending_overlay_write.mode, ) return res