mirror of
https://github.com/eigent-ai/eigent.git
synced 2026-08-27 17:41:56 +00:00
361 lines
13 KiB
Python
361 lines
13 KiB
Python
# ========= 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
|