mirror of
https://github.com/langgenius/dify.git
synced 2026-09-24 23:22:26 +08:00
feat(agent): Sandbox / CLI Agent (dify.shell) + read-only sandbox file inspector (#36984)
Co-authored-by: Claude Opus 4.8 (1M context) <noreply@anthropic.com> Co-authored-by: autofix-ci[bot] <114827586+autofix-ci[bot]@users.noreply.github.com>
This commit is contained in:
co-authored by
Claude Opus 4.8
autofix-ci[bot]
parent
d3058d63bd
commit
44725dde74
@@ -30,6 +30,7 @@ from dify_agent.layers.execution_context import (
|
||||
DifyExecutionContextLayerConfig,
|
||||
)
|
||||
from dify_agent.layers.output import DIFY_OUTPUT_LAYER_TYPE_ID, DifyOutputLayerConfig
|
||||
from dify_agent.layers.shell import DIFY_SHELL_LAYER_TYPE_ID, DifyShellLayerConfig
|
||||
from dify_agent.protocol import (
|
||||
DIFY_AGENT_HISTORY_LAYER_ID,
|
||||
DIFY_AGENT_MODEL_LAYER_ID,
|
||||
@@ -48,6 +49,7 @@ WORKFLOW_USER_PROMPT_LAYER_ID = "workflow_user_prompt"
|
||||
AGENT_APP_USER_PROMPT_LAYER_ID = "agent_app_user_prompt"
|
||||
DIFY_EXECUTION_CONTEXT_LAYER_ID = "execution_context"
|
||||
DIFY_PLUGIN_TOOLS_LAYER_ID = "tools"
|
||||
DIFY_SHELL_LAYER_ID = "shell"
|
||||
|
||||
# Layer types that hold credentials in their per-run config. These are excluded
|
||||
# from the cleanup-replay composition (and from the snapshot that is sent with
|
||||
@@ -167,6 +169,10 @@ class AgentBackendWorkflowNodeRunInput(BaseModel):
|
||||
idempotency_key: str | None = None
|
||||
output: AgentBackendOutputConfig | None = None
|
||||
tools: DifyPluginToolsLayerConfig | None = None
|
||||
# Inject the sandboxed shell layer (dify.shell). Requires the agent backend
|
||||
# to be wired with a shellctl entrypoint; see configs AGENT_SHELL_ENABLED.
|
||||
include_shell: bool = False
|
||||
shell_config: DifyShellLayerConfig | None = None
|
||||
session_snapshot: CompositorSessionSnapshot | None = None
|
||||
include_history: bool = True
|
||||
suspend_on_exit: bool = True
|
||||
@@ -199,6 +205,10 @@ class AgentBackendAgentAppRunInput(BaseModel):
|
||||
idempotency_key: str | None = None
|
||||
output: AgentBackendOutputConfig | None = None
|
||||
tools: DifyPluginToolsLayerConfig | None = None
|
||||
# Inject the sandboxed shell layer (dify.shell). Requires the agent backend
|
||||
# to be wired with a shellctl entrypoint; see configs AGENT_SHELL_ENABLED.
|
||||
include_shell: bool = False
|
||||
shell_config: DifyShellLayerConfig | None = None
|
||||
session_snapshot: CompositorSessionSnapshot | None = None
|
||||
include_history: bool = True
|
||||
suspend_on_exit: bool = True
|
||||
@@ -289,6 +299,18 @@ class AgentBackendRunRequestBuilder:
|
||||
)
|
||||
)
|
||||
|
||||
if run_input.include_shell:
|
||||
# Sandboxed bash workspace (dify.shell). The layer declares NoLayerDeps,
|
||||
# so the spec carries no deps; shellctl connection is server-injected.
|
||||
layers.append(
|
||||
RunLayerSpec(
|
||||
name=DIFY_SHELL_LAYER_ID,
|
||||
type=DIFY_SHELL_LAYER_TYPE_ID,
|
||||
metadata=run_input.metadata,
|
||||
config=run_input.shell_config or DifyShellLayerConfig(),
|
||||
)
|
||||
)
|
||||
|
||||
if run_input.output is not None:
|
||||
layers.append(
|
||||
RunLayerSpec(
|
||||
@@ -432,6 +454,18 @@ class AgentBackendRunRequestBuilder:
|
||||
)
|
||||
)
|
||||
|
||||
if run_input.include_shell:
|
||||
# Sandboxed bash workspace (dify.shell). The layer declares NoLayerDeps,
|
||||
# so the spec carries no deps; shellctl connection is server-injected.
|
||||
layers.append(
|
||||
RunLayerSpec(
|
||||
name=DIFY_SHELL_LAYER_ID,
|
||||
type=DIFY_SHELL_LAYER_TYPE_ID,
|
||||
metadata=run_input.metadata,
|
||||
config=run_input.shell_config or DifyShellLayerConfig(),
|
||||
)
|
||||
)
|
||||
|
||||
if run_input.output is not None:
|
||||
layers.append(
|
||||
RunLayerSpec(
|
||||
|
||||
@@ -0,0 +1,135 @@
|
||||
"""API-side client for the agent backend's read-only workspace file endpoints.
|
||||
|
||||
The agent backend exposes ``/workspaces/{session_id}/files{,/preview,/download}``
|
||||
to inspect a shell-layer sandbox workspace. This thin synchronous client proxies
|
||||
those reads for the console FS inspector and normalizes transport/HTTP failures
|
||||
into the API backend's ``AgentBackendError`` boundary, preserving the backend's
|
||||
status code and ``{code, message}`` detail so the controller can relay them.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import base64
|
||||
import binascii
|
||||
from dataclasses import dataclass
|
||||
from typing import Literal
|
||||
|
||||
import httpx
|
||||
from pydantic import BaseModel
|
||||
|
||||
from clients.agent_backend.errors import AgentBackendHTTPError, AgentBackendTransportError
|
||||
|
||||
_DEFAULT_TIMEOUT_SECONDS = 30.0
|
||||
|
||||
|
||||
class WorkspaceFileEntry(BaseModel):
|
||||
"""One entry in a workspace directory listing."""
|
||||
|
||||
name: str
|
||||
type: Literal["file", "dir", "symlink"]
|
||||
size: int
|
||||
mtime: int
|
||||
|
||||
|
||||
class WorkspaceListResult(BaseModel):
|
||||
"""Directory listing of a workspace path."""
|
||||
|
||||
path: str
|
||||
entries: list[WorkspaceFileEntry]
|
||||
truncated: bool
|
||||
|
||||
|
||||
class WorkspacePreviewResult(BaseModel):
|
||||
"""Inline preview of a workspace file."""
|
||||
|
||||
path: str
|
||||
size: int
|
||||
truncated: bool
|
||||
binary: bool
|
||||
text: str | None = None
|
||||
|
||||
|
||||
@dataclass(frozen=True, slots=True)
|
||||
class WorkspaceDownloadResult:
|
||||
"""Decoded bytes of a workspace file for download."""
|
||||
|
||||
path: str
|
||||
size: int
|
||||
truncated: bool
|
||||
content: bytes
|
||||
|
||||
|
||||
class WorkspaceFilesBackendClient:
|
||||
"""Synchronous proxy to the agent backend workspace file endpoints."""
|
||||
|
||||
def __init__(
|
||||
self,
|
||||
base_url: str,
|
||||
*,
|
||||
timeout: float = _DEFAULT_TIMEOUT_SECONDS,
|
||||
transport: httpx.BaseTransport | None = None,
|
||||
) -> None:
|
||||
self._base_url = base_url.rstrip("/")
|
||||
self._timeout = timeout
|
||||
self._transport = transport
|
||||
|
||||
def list_files(self, session_id: str, path: str) -> WorkspaceListResult:
|
||||
data = self._get(f"/workspaces/{session_id}/files", params={"path": path})
|
||||
return WorkspaceListResult.model_validate(data)
|
||||
|
||||
def preview(self, session_id: str, path: str) -> WorkspacePreviewResult:
|
||||
data = self._get(f"/workspaces/{session_id}/files/preview", params={"path": path})
|
||||
return WorkspacePreviewResult.model_validate(data)
|
||||
|
||||
def download(self, session_id: str, path: str) -> WorkspaceDownloadResult:
|
||||
data = self._get(f"/workspaces/{session_id}/files/download", params={"path": path})
|
||||
encoded = data.get("content_base64")
|
||||
if not isinstance(encoded, str):
|
||||
raise AgentBackendHTTPError("agent backend download response missing content", status_code=502, detail=data)
|
||||
try:
|
||||
content = base64.b64decode(encoded, validate=True)
|
||||
except (binascii.Error, ValueError) as exc:
|
||||
raise AgentBackendHTTPError(
|
||||
"agent backend returned undecodable download content", status_code=502, detail=str(exc)
|
||||
) from exc
|
||||
size = data.get("size")
|
||||
return WorkspaceDownloadResult(
|
||||
path=str(data.get("path", path)),
|
||||
size=int(size) if isinstance(size, (int, float)) else len(content),
|
||||
truncated=bool(data.get("truncated")),
|
||||
content=content,
|
||||
)
|
||||
|
||||
def _get(self, route: str, *, params: dict[str, str]) -> dict[str, object]:
|
||||
url = f"{self._base_url}{route}"
|
||||
try:
|
||||
with httpx.Client(timeout=self._timeout, transport=self._transport, trust_env=False) as client:
|
||||
response = client.get(url, params=params)
|
||||
except httpx.HTTPError as exc:
|
||||
raise AgentBackendTransportError(f"failed to reach agent backend workspace endpoint: {exc}") from exc
|
||||
if response.status_code >= 400:
|
||||
detail: object
|
||||
try:
|
||||
detail = response.json().get("detail", response.text)
|
||||
except ValueError:
|
||||
detail = response.text
|
||||
raise AgentBackendHTTPError(
|
||||
f"agent backend workspace request failed ({response.status_code})",
|
||||
status_code=response.status_code,
|
||||
detail=detail,
|
||||
)
|
||||
body = response.json()
|
||||
if not isinstance(body, dict):
|
||||
raise AgentBackendHTTPError(
|
||||
"agent backend workspace response was not an object", status_code=502, detail=body
|
||||
)
|
||||
return body
|
||||
|
||||
|
||||
__all__ = [
|
||||
"WorkspaceDownloadResult",
|
||||
"WorkspaceFileEntry",
|
||||
"WorkspaceFilesBackendClient",
|
||||
"WorkspaceListResult",
|
||||
"WorkspacePreviewResult",
|
||||
]
|
||||
@@ -21,3 +21,13 @@ class AgentBackendConfig(BaseSettings):
|
||||
description="Scenario used by the fake Agent backend client.",
|
||||
default="success",
|
||||
)
|
||||
|
||||
AGENT_SHELL_ENABLED: bool = Field(
|
||||
description=(
|
||||
"Inject the dify.shell layer (sandboxed bash workspace) into Agent runs. "
|
||||
"Requires the agent backend to be wired with a shellctl entrypoint; keep it "
|
||||
"off until shellctl is deployed, otherwise every agent run that includes the "
|
||||
"shell layer will fail."
|
||||
),
|
||||
default=False,
|
||||
)
|
||||
|
||||
@@ -53,6 +53,7 @@ from .app import (
|
||||
agent,
|
||||
agent_app_access,
|
||||
agent_app_feature,
|
||||
agent_app_workspace,
|
||||
annotation,
|
||||
app,
|
||||
audio,
|
||||
@@ -150,6 +151,7 @@ __all__ = [
|
||||
"agent",
|
||||
"agent_app_access",
|
||||
"agent_app_feature",
|
||||
"agent_app_workspace",
|
||||
"agent_composer",
|
||||
"agent_providers",
|
||||
"agent_roster",
|
||||
|
||||
@@ -0,0 +1,319 @@
|
||||
"""Agent App sandbox file-system inspector (read-only).
|
||||
|
||||
Exposes the PRD "rc1-like sandbox file system, downloadable not editable" view
|
||||
for an Agent App conversation: list a directory, preview a file, or download a
|
||||
file from the conversation's shell-layer workspace. The API never touches
|
||||
shellctl directly — it resolves the conversation's sandbox ``session_id`` from
|
||||
the stored session snapshot and proxies to the agent backend's read-only
|
||||
workspace endpoints.
|
||||
"""
|
||||
|
||||
from typing import Literal
|
||||
from uuid import UUID
|
||||
|
||||
from flask import Response
|
||||
from flask_restx import Resource, fields
|
||||
from pydantic import BaseModel, Field
|
||||
|
||||
from clients.agent_backend.errors import AgentBackendHTTPError, AgentBackendTransportError
|
||||
from clients.agent_backend.workspace_files_client import WorkspaceDownloadResult
|
||||
from controllers.common.schema import (
|
||||
query_params_from_model,
|
||||
query_params_from_request,
|
||||
register_response_schema_models,
|
||||
)
|
||||
from controllers.console import console_ns
|
||||
from controllers.console.app.wraps import get_app_model
|
||||
from controllers.console.wraps import account_initialization_required, setup_required
|
||||
from fields.base import ResponseModel
|
||||
from libs.login import current_account_with_tenant, login_required
|
||||
from models.model import App, AppMode
|
||||
from services.agent_app_workspace_service import (
|
||||
AgentAppWorkspaceService,
|
||||
AgentWorkspaceInspectorError,
|
||||
WorkflowAgentWorkspaceService,
|
||||
)
|
||||
|
||||
|
||||
class _WorkspaceFileDownloadField(fields.Raw):
|
||||
__schema_type__ = "string"
|
||||
__schema_format__ = "binary"
|
||||
|
||||
|
||||
class AgentWorkspaceListQuery(BaseModel):
|
||||
conversation_id: str = Field(min_length=1, description="Agent App conversation ID")
|
||||
path: str = Field(default=".", description="Directory path relative to the sandbox workspace")
|
||||
|
||||
|
||||
class AgentWorkspaceFileQuery(BaseModel):
|
||||
conversation_id: str = Field(min_length=1, description="Agent App conversation ID")
|
||||
path: str = Field(min_length=1, description="File path relative to the sandbox workspace")
|
||||
|
||||
|
||||
class WorkflowAgentWorkspaceListQuery(BaseModel):
|
||||
path: str = Field(default=".", description="Directory path relative to the sandbox workspace")
|
||||
node_execution_id: str | None = Field(
|
||||
default=None,
|
||||
description=(
|
||||
"Optional workflow node execution ID. When omitted, the latest active session for the node is used."
|
||||
),
|
||||
)
|
||||
|
||||
|
||||
class WorkflowAgentWorkspaceFileQuery(BaseModel):
|
||||
path: str = Field(min_length=1, description="File path relative to the sandbox workspace")
|
||||
node_execution_id: str | None = Field(
|
||||
default=None,
|
||||
description=(
|
||||
"Optional workflow node execution ID. When omitted, the latest active session for the node is used."
|
||||
),
|
||||
)
|
||||
|
||||
|
||||
class WorkspaceFileEntryResponse(ResponseModel):
|
||||
name: str
|
||||
type: Literal["file", "dir", "symlink"]
|
||||
size: int
|
||||
mtime: int
|
||||
|
||||
|
||||
class WorkspaceListResponse(ResponseModel):
|
||||
path: str
|
||||
entries: list[WorkspaceFileEntryResponse] = Field(default_factory=list)
|
||||
truncated: bool = False
|
||||
|
||||
|
||||
class WorkspacePreviewResponse(ResponseModel):
|
||||
path: str
|
||||
size: int
|
||||
truncated: bool
|
||||
binary: bool
|
||||
text: str | None = None
|
||||
|
||||
|
||||
register_response_schema_models(console_ns, WorkspaceListResponse)
|
||||
register_response_schema_models(console_ns, WorkspacePreviewResponse)
|
||||
|
||||
|
||||
def _handle(exc: Exception) -> tuple[dict[str, object], int]:
|
||||
if isinstance(exc, AgentWorkspaceInspectorError):
|
||||
return {"code": exc.code, "message": exc.message}, exc.status_code
|
||||
if isinstance(exc, AgentBackendHTTPError):
|
||||
detail = exc.detail
|
||||
if isinstance(detail, dict):
|
||||
return {
|
||||
"code": detail.get("code", "agent_backend_error"),
|
||||
"message": detail.get("message", str(exc)),
|
||||
}, exc.status_code
|
||||
return {"code": "agent_backend_error", "message": str(detail)}, exc.status_code
|
||||
if isinstance(exc, AgentBackendTransportError):
|
||||
return {"code": "agent_backend_unreachable", "message": str(exc)}, 502
|
||||
raise exc
|
||||
|
||||
|
||||
def _download_response(result: WorkspaceDownloadResult) -> Response | tuple[dict[str, object], int]:
|
||||
if result.truncated:
|
||||
return {
|
||||
"code": "workspace_file_too_large",
|
||||
"message": (
|
||||
"file exceeds the workspace download limit; use preview for partial text or download a smaller file"
|
||||
),
|
||||
"size": result.size,
|
||||
}, 413
|
||||
filename = result.path.rsplit("/", 1)[-1] or "download"
|
||||
return Response(
|
||||
result.content,
|
||||
mimetype="application/octet-stream",
|
||||
headers={
|
||||
"Content-Disposition": f'attachment; filename="{filename}"',
|
||||
"Content-Length": str(len(result.content)),
|
||||
"X-Workspace-File-Size": str(result.size),
|
||||
},
|
||||
)
|
||||
|
||||
|
||||
@console_ns.route("/apps/<uuid:app_id>/agent-workspace/files")
|
||||
class AgentAppWorkspaceListResource(Resource):
|
||||
@console_ns.doc("list_agent_app_workspace_files")
|
||||
@console_ns.doc(description="List a directory in an Agent App conversation's sandbox workspace (read-only)")
|
||||
@console_ns.doc(params={"app_id": "Application ID", **query_params_from_model(AgentWorkspaceListQuery)})
|
||||
@console_ns.response(200, "Listing returned", console_ns.models[WorkspaceListResponse.__name__])
|
||||
@setup_required
|
||||
@login_required
|
||||
@account_initialization_required
|
||||
@get_app_model(mode=[AppMode.AGENT])
|
||||
def get(self, app_model: App):
|
||||
_, tenant_id = current_account_with_tenant()
|
||||
query = query_params_from_request(AgentWorkspaceListQuery)
|
||||
try:
|
||||
result = AgentAppWorkspaceService().list_files(
|
||||
tenant_id=tenant_id,
|
||||
app_id=app_model.id,
|
||||
conversation_id=query.conversation_id,
|
||||
path=query.path,
|
||||
)
|
||||
except Exception as exc: # normalized to an HTTP response below
|
||||
return _handle(exc)
|
||||
return result.model_dump()
|
||||
|
||||
|
||||
@console_ns.route("/apps/<uuid:app_id>/agent-workspace/files/preview")
|
||||
class AgentAppWorkspacePreviewResource(Resource):
|
||||
@console_ns.doc("preview_agent_app_workspace_file")
|
||||
@console_ns.doc(description="Preview a text/binary file in an Agent App conversation's sandbox workspace")
|
||||
@console_ns.doc(params={"app_id": "Application ID", **query_params_from_model(AgentWorkspaceFileQuery)})
|
||||
@console_ns.response(200, "Preview returned", console_ns.models[WorkspacePreviewResponse.__name__])
|
||||
@setup_required
|
||||
@login_required
|
||||
@account_initialization_required
|
||||
@get_app_model(mode=[AppMode.AGENT])
|
||||
def get(self, app_model: App):
|
||||
_, tenant_id = current_account_with_tenant()
|
||||
query = query_params_from_request(AgentWorkspaceFileQuery)
|
||||
try:
|
||||
result = AgentAppWorkspaceService().preview(
|
||||
tenant_id=tenant_id,
|
||||
app_id=app_model.id,
|
||||
conversation_id=query.conversation_id,
|
||||
path=query.path,
|
||||
)
|
||||
except Exception as exc: # normalized to an HTTP response below
|
||||
return _handle(exc)
|
||||
return result.model_dump()
|
||||
|
||||
|
||||
@console_ns.route("/apps/<uuid:app_id>/agent-workspace/files/download")
|
||||
class AgentAppWorkspaceDownloadResource(Resource):
|
||||
@console_ns.doc("download_agent_app_workspace_file")
|
||||
@console_ns.doc(description="Download a file from an Agent App conversation's sandbox workspace (read-only)")
|
||||
@console_ns.doc(params={"app_id": "Application ID", **query_params_from_model(AgentWorkspaceFileQuery)})
|
||||
@console_ns.doc(produces=["application/octet-stream"])
|
||||
@console_ns.response(200, "File bytes", _WorkspaceFileDownloadField)
|
||||
@console_ns.response(413, "File exceeds the workspace download limit")
|
||||
@setup_required
|
||||
@login_required
|
||||
@account_initialization_required
|
||||
@get_app_model(mode=[AppMode.AGENT])
|
||||
def get(self, app_model: App):
|
||||
_, tenant_id = current_account_with_tenant()
|
||||
query = query_params_from_request(AgentWorkspaceFileQuery)
|
||||
try:
|
||||
result = AgentAppWorkspaceService().download(
|
||||
tenant_id=tenant_id,
|
||||
app_id=app_model.id,
|
||||
conversation_id=query.conversation_id,
|
||||
path=query.path,
|
||||
)
|
||||
except Exception as exc: # normalized to an HTTP response below
|
||||
return _handle(exc)
|
||||
return _download_response(result)
|
||||
|
||||
|
||||
@console_ns.route(
|
||||
"/apps/<uuid:app_id>/workflow-runs/<uuid:workflow_run_id>/agent-nodes/<string:node_id>/workspace/files"
|
||||
)
|
||||
class WorkflowAgentWorkspaceListResource(Resource):
|
||||
@console_ns.doc("list_workflow_agent_workspace_files")
|
||||
@console_ns.doc(description="List a directory in a Workflow Agent node's sandbox workspace (read-only)")
|
||||
@console_ns.doc(
|
||||
params={
|
||||
"app_id": "Application ID",
|
||||
"workflow_run_id": "Workflow run ID",
|
||||
"node_id": "Workflow Agent node ID",
|
||||
**query_params_from_model(WorkflowAgentWorkspaceListQuery),
|
||||
}
|
||||
)
|
||||
@console_ns.response(200, "Listing returned", console_ns.models[WorkspaceListResponse.__name__])
|
||||
@setup_required
|
||||
@login_required
|
||||
@account_initialization_required
|
||||
@get_app_model(mode=[AppMode.ADVANCED_CHAT, AppMode.WORKFLOW])
|
||||
def get(self, app_model: App, workflow_run_id: UUID, node_id: str):
|
||||
_, tenant_id = current_account_with_tenant()
|
||||
query = query_params_from_request(WorkflowAgentWorkspaceListQuery)
|
||||
try:
|
||||
result = WorkflowAgentWorkspaceService().list_files(
|
||||
tenant_id=tenant_id,
|
||||
app_id=app_model.id,
|
||||
workflow_run_id=str(workflow_run_id),
|
||||
node_id=node_id,
|
||||
node_execution_id=query.node_execution_id,
|
||||
path=query.path,
|
||||
)
|
||||
except Exception as exc: # normalized to an HTTP response below
|
||||
return _handle(exc)
|
||||
return result.model_dump()
|
||||
|
||||
|
||||
@console_ns.route(
|
||||
"/apps/<uuid:app_id>/workflow-runs/<uuid:workflow_run_id>/agent-nodes/<string:node_id>/workspace/files/preview"
|
||||
)
|
||||
class WorkflowAgentWorkspacePreviewResource(Resource):
|
||||
@console_ns.doc("preview_workflow_agent_workspace_file")
|
||||
@console_ns.doc(description="Preview a text/binary file in a Workflow Agent node's sandbox workspace")
|
||||
@console_ns.doc(
|
||||
params={
|
||||
"app_id": "Application ID",
|
||||
"workflow_run_id": "Workflow run ID",
|
||||
"node_id": "Workflow Agent node ID",
|
||||
**query_params_from_model(WorkflowAgentWorkspaceFileQuery),
|
||||
}
|
||||
)
|
||||
@console_ns.response(200, "Preview returned", console_ns.models[WorkspacePreviewResponse.__name__])
|
||||
@setup_required
|
||||
@login_required
|
||||
@account_initialization_required
|
||||
@get_app_model(mode=[AppMode.ADVANCED_CHAT, AppMode.WORKFLOW])
|
||||
def get(self, app_model: App, workflow_run_id: UUID, node_id: str):
|
||||
_, tenant_id = current_account_with_tenant()
|
||||
query = query_params_from_request(WorkflowAgentWorkspaceFileQuery)
|
||||
try:
|
||||
result = WorkflowAgentWorkspaceService().preview(
|
||||
tenant_id=tenant_id,
|
||||
app_id=app_model.id,
|
||||
workflow_run_id=str(workflow_run_id),
|
||||
node_id=node_id,
|
||||
node_execution_id=query.node_execution_id,
|
||||
path=query.path,
|
||||
)
|
||||
except Exception as exc: # normalized to an HTTP response below
|
||||
return _handle(exc)
|
||||
return result.model_dump()
|
||||
|
||||
|
||||
@console_ns.route(
|
||||
"/apps/<uuid:app_id>/workflow-runs/<uuid:workflow_run_id>/agent-nodes/<string:node_id>/workspace/files/download"
|
||||
)
|
||||
class WorkflowAgentWorkspaceDownloadResource(Resource):
|
||||
@console_ns.doc("download_workflow_agent_workspace_file")
|
||||
@console_ns.doc(description="Download a file from a Workflow Agent node's sandbox workspace (read-only)")
|
||||
@console_ns.doc(
|
||||
params={
|
||||
"app_id": "Application ID",
|
||||
"workflow_run_id": "Workflow run ID",
|
||||
"node_id": "Workflow Agent node ID",
|
||||
**query_params_from_model(WorkflowAgentWorkspaceFileQuery),
|
||||
}
|
||||
)
|
||||
@console_ns.doc(produces=["application/octet-stream"])
|
||||
@console_ns.response(200, "File bytes", _WorkspaceFileDownloadField)
|
||||
@console_ns.response(413, "File exceeds the workspace download limit")
|
||||
@setup_required
|
||||
@login_required
|
||||
@account_initialization_required
|
||||
@get_app_model(mode=[AppMode.ADVANCED_CHAT, AppMode.WORKFLOW])
|
||||
def get(self, app_model: App, workflow_run_id: UUID, node_id: str):
|
||||
_, tenant_id = current_account_with_tenant()
|
||||
query = query_params_from_request(WorkflowAgentWorkspaceFileQuery)
|
||||
try:
|
||||
result = WorkflowAgentWorkspaceService().download(
|
||||
tenant_id=tenant_id,
|
||||
app_id=app_model.id,
|
||||
workflow_run_id=str(workflow_run_id),
|
||||
node_id=node_id,
|
||||
node_execution_id=query.node_execution_id,
|
||||
path=query.path,
|
||||
)
|
||||
except Exception as exc: # normalized to an HTTP response below
|
||||
return _handle(exc)
|
||||
return _download_response(result)
|
||||
@@ -22,11 +22,13 @@ from clients.agent_backend import (
|
||||
AgentBackendRunRequestBuilder,
|
||||
redact_for_agent_backend_log,
|
||||
)
|
||||
from configs import dify_config
|
||||
from core.app.entities.app_invoke_entities import DifyRunContext
|
||||
from core.workflow.nodes.agent_v2.plugin_tools_builder import (
|
||||
WorkflowAgentPluginToolsBuilder,
|
||||
WorkflowAgentPluginToolsBuildError,
|
||||
)
|
||||
from core.workflow.nodes.agent_v2.runtime_request_builder import build_shell_layer_config
|
||||
from models.agent_config_entities import AgentSoulConfig
|
||||
from models.provider_ids import ModelProviderID
|
||||
|
||||
@@ -96,10 +98,13 @@ class AgentAppRuntimeRequestBuilder:
|
||||
)
|
||||
except WorkflowAgentPluginToolsBuildError as error:
|
||||
raise AgentAppRuntimeRequestBuildError(error.error_code, str(error)) from error
|
||||
if tools_layer is not None:
|
||||
if tools_layer is not None or agent_soul.tools.cli_tools:
|
||||
metadata["agent_tools"] = {
|
||||
"dify_tool_count": len(tools_layer.tools),
|
||||
"dify_tool_names": [tool.name or tool.tool_name for tool in tools_layer.tools],
|
||||
"dify_tool_count": len(tools_layer.tools) if tools_layer is not None else 0,
|
||||
"dify_tool_names": [tool.name or tool.tool_name for tool in tools_layer.tools]
|
||||
if tools_layer is not None
|
||||
else [],
|
||||
"cli_tool_count": len(agent_soul.tools.cli_tools),
|
||||
}
|
||||
|
||||
request = self._request_builder.build_for_agent_app(
|
||||
@@ -126,6 +131,8 @@ class AgentAppRuntimeRequestBuilder:
|
||||
agent_soul_prompt=agent_soul.prompt.system_prompt or None,
|
||||
user_prompt=context.user_query,
|
||||
tools=tools_layer,
|
||||
include_shell=dify_config.AGENT_SHELL_ENABLED,
|
||||
shell_config=build_shell_layer_config(agent_soul),
|
||||
session_snapshot=context.session_snapshot,
|
||||
idempotency_key=context.idempotency_key,
|
||||
metadata=metadata,
|
||||
|
||||
@@ -44,6 +44,32 @@ class AgentAppRuntimeSessionStore:
|
||||
return None
|
||||
return CompositorSessionSnapshot.model_validate_json(row.session_snapshot)
|
||||
|
||||
def load_active_snapshot_for_conversation(
|
||||
self, *, tenant_id: str, app_id: str, conversation_id: str
|
||||
) -> CompositorSessionSnapshot | None:
|
||||
"""Load a conversation's active snapshot without the agent/config scope.
|
||||
|
||||
One Agent App conversation maps to one active session, so the workspace
|
||||
inspector can resolve it from the conversation alone (it does not know
|
||||
which agent config version a past turn ran under).
|
||||
"""
|
||||
stmt = (
|
||||
select(AgentRuntimeSession)
|
||||
.where(
|
||||
AgentRuntimeSession.owner_type == AgentRuntimeSessionOwnerType.CONVERSATION,
|
||||
AgentRuntimeSession.tenant_id == tenant_id,
|
||||
AgentRuntimeSession.app_id == app_id,
|
||||
AgentRuntimeSession.conversation_id == conversation_id,
|
||||
AgentRuntimeSession.status == AgentRuntimeSessionStatus.ACTIVE,
|
||||
)
|
||||
.order_by(AgentRuntimeSession.updated_at.desc())
|
||||
)
|
||||
with session_factory.create_session() as session:
|
||||
row = session.scalar(stmt)
|
||||
if row is None:
|
||||
return None
|
||||
return CompositorSessionSnapshot.model_validate_json(row.session_snapshot)
|
||||
|
||||
def save_active_snapshot(
|
||||
self,
|
||||
*,
|
||||
@@ -75,6 +101,20 @@ class AgentAppRuntimeSessionStore:
|
||||
row.session_snapshot = snapshot_json
|
||||
row.status = AgentRuntimeSessionStatus.ACTIVE
|
||||
row.cleaned_at = None
|
||||
session.flush()
|
||||
other_rows = session.scalars(
|
||||
select(AgentRuntimeSession).where(
|
||||
AgentRuntimeSession.owner_type == AgentRuntimeSessionOwnerType.CONVERSATION,
|
||||
AgentRuntimeSession.tenant_id == scope.tenant_id,
|
||||
AgentRuntimeSession.app_id == scope.app_id,
|
||||
AgentRuntimeSession.conversation_id == scope.conversation_id,
|
||||
AgentRuntimeSession.status == AgentRuntimeSessionStatus.ACTIVE,
|
||||
AgentRuntimeSession.id != row.id,
|
||||
)
|
||||
).all()
|
||||
for other_row in other_rows:
|
||||
other_row.status = AgentRuntimeSessionStatus.CLEANED
|
||||
other_row.cleaned_at = naive_utc_now()
|
||||
session.commit()
|
||||
|
||||
def mark_cleaned(self, *, scope: AgentAppSessionScope, backend_run_id: str | None = None) -> None:
|
||||
|
||||
@@ -12,24 +12,24 @@ SUPPORTED_AGENT_BACKEND_FEATURES = frozenset(
|
||||
"model",
|
||||
"structured_output",
|
||||
"tools.dify_tools",
|
||||
"tools.cli_tools",
|
||||
"env",
|
||||
"sandbox",
|
||||
}
|
||||
)
|
||||
|
||||
RESERVED_AGENT_BACKEND_FEATURES = frozenset(
|
||||
{
|
||||
"skills_files",
|
||||
"tools.cli_tools",
|
||||
"knowledge",
|
||||
"human",
|
||||
"env",
|
||||
"sandbox",
|
||||
"memory",
|
||||
}
|
||||
)
|
||||
|
||||
|
||||
def build_runtime_feature_manifest(agent_soul: AgentSoulConfig) -> dict[str, Any]:
|
||||
"""Describe PRD capabilities that are persisted but not executed in phase 3."""
|
||||
"""Describe PRD capabilities supported by or still reserved from Agent backend runtime."""
|
||||
warnings: list[dict[str, str]] = []
|
||||
soul_dump = agent_soul.model_dump(mode="json", exclude_none=True, exclude_defaults=True)
|
||||
for section in sorted(RESERVED_AGENT_BACKEND_FEATURES):
|
||||
@@ -48,6 +48,9 @@ def build_runtime_feature_manifest(agent_soul: AgentSoulConfig) -> dict[str, Any
|
||||
|
||||
reserved_status = dict.fromkeys(sorted(RESERVED_AGENT_BACKEND_FEATURES), "reserved_not_executed")
|
||||
reserved_status["tools.dify_tools"] = "supported_when_config_valid"
|
||||
reserved_status["tools.cli_tools"] = "supported_by_shell_bootstrap"
|
||||
reserved_status["env"] = "supported_by_shell_bootstrap"
|
||||
reserved_status["sandbox"] = "forwarded_to_shell_layer_config"
|
||||
|
||||
return {
|
||||
"supported": sorted(SUPPORTED_AGENT_BACKEND_FEATURES),
|
||||
|
||||
@@ -2,11 +2,19 @@ from __future__ import annotations
|
||||
|
||||
from collections.abc import Mapping, Sequence
|
||||
from dataclasses import dataclass
|
||||
from typing import Any, Literal, Protocol, cast
|
||||
from typing import Any, Literal, Protocol, assert_never, cast
|
||||
|
||||
from agenton.compositor import CompositorSessionSnapshot
|
||||
from dify_agent.layers.execution_context import DifyExecutionContextLayerConfig
|
||||
from dify_agent.layers.shell import (
|
||||
DifyShellCliToolConfig,
|
||||
DifyShellEnvVarConfig,
|
||||
DifyShellLayerConfig,
|
||||
DifyShellSandboxConfig,
|
||||
DifyShellSecretRefConfig,
|
||||
)
|
||||
from dify_agent.protocol import CreateRunRequest
|
||||
from pydantic import BaseModel
|
||||
|
||||
from clients.agent_backend import (
|
||||
AgentBackendModelConfig,
|
||||
@@ -15,6 +23,7 @@ from clients.agent_backend import (
|
||||
AgentBackendWorkflowNodeRunInput,
|
||||
redact_for_agent_backend_log,
|
||||
)
|
||||
from configs import dify_config
|
||||
from core.app.entities.app_invoke_entities import DifyRunContext, InvokeFrom
|
||||
from core.workflow.system_variables import SystemVariableKey, get_system_text
|
||||
from graphon.variables.segments import Segment
|
||||
@@ -123,10 +132,12 @@ class WorkflowAgentRuntimeRequestBuilder:
|
||||
)
|
||||
except WorkflowAgentPluginToolsBuildError as error:
|
||||
raise WorkflowAgentRuntimeRequestBuildError(error.error_code, str(error)) from error
|
||||
if tools_layer is not None:
|
||||
if tools_layer is not None or agent_soul.tools.cli_tools:
|
||||
metadata["agent_tools"] = {
|
||||
"dify_tool_count": len(tools_layer.tools),
|
||||
"dify_tool_names": [tool.name or tool.tool_name for tool in tools_layer.tools],
|
||||
"dify_tool_count": len(tools_layer.tools) if tools_layer is not None else 0,
|
||||
"dify_tool_names": [tool.name or tool.tool_name for tool in tools_layer.tools]
|
||||
if tools_layer is not None
|
||||
else [],
|
||||
"cli_tool_count": len(agent_soul.tools.cli_tools),
|
||||
}
|
||||
|
||||
@@ -165,6 +176,8 @@ class WorkflowAgentRuntimeRequestBuilder:
|
||||
user_prompt=user_prompt,
|
||||
output=self._build_output_config(node_job.declared_outputs),
|
||||
tools=tools_layer,
|
||||
include_shell=dify_config.AGENT_SHELL_ENABLED,
|
||||
shell_config=build_shell_layer_config(agent_soul),
|
||||
session_snapshot=context.session_snapshot,
|
||||
idempotency_key=self._idempotency_key(context),
|
||||
metadata=metadata,
|
||||
@@ -372,7 +385,9 @@ class WorkflowAgentRuntimeRequestBuilder:
|
||||
"mime_type": {"type": "string"},
|
||||
"url": {"type": "string"},
|
||||
},
|
||||
"required": ["file_id"],
|
||||
}
|
||||
assert_never(output_type)
|
||||
|
||||
@staticmethod
|
||||
def _normalize_credentials(credentials: Mapping[str, Any]) -> dict[str, str | int | float | bool | None]:
|
||||
@@ -383,3 +398,73 @@ class WorkflowAgentRuntimeRequestBuilder:
|
||||
else:
|
||||
normalized[key] = str(value)
|
||||
return normalized
|
||||
|
||||
|
||||
def build_shell_layer_config(agent_soul: AgentSoulConfig) -> DifyShellLayerConfig:
|
||||
"""Map Agent Soul shell-adjacent fields into the Agent backend shell config."""
|
||||
sandbox_config = _plain_mapping(agent_soul.sandbox.config)
|
||||
return DifyShellLayerConfig(
|
||||
cli_tools=[tool for tool in (_shell_cli_tool(item) for item in agent_soul.tools.cli_tools) if tool is not None],
|
||||
env=[env for env in (_shell_env_var(item) for item in agent_soul.env.variables) if env is not None],
|
||||
secret_refs=[
|
||||
secret for secret in (_shell_secret_ref(item) for item in agent_soul.env.secret_refs) if secret is not None
|
||||
],
|
||||
sandbox=DifyShellSandboxConfig(
|
||||
provider=agent_soul.sandbox.provider,
|
||||
config=sandbox_config,
|
||||
)
|
||||
if agent_soul.sandbox.provider or sandbox_config
|
||||
else None,
|
||||
)
|
||||
|
||||
|
||||
def _shell_cli_tool(item: object) -> DifyShellCliToolConfig | None:
|
||||
data = _plain_mapping(item)
|
||||
commands: list[str] = []
|
||||
raw_commands = data.get("install_commands")
|
||||
if isinstance(raw_commands, list):
|
||||
commands.extend(str(command) for command in raw_commands if str(command).strip())
|
||||
for key in ("install_command", "install", "setup_command"):
|
||||
raw_command = data.get(key)
|
||||
if isinstance(raw_command, str) and raw_command.strip():
|
||||
commands.append(raw_command)
|
||||
name = data.get("name") or data.get("tool_name") or data.get("label")
|
||||
if not commands and not isinstance(name, str):
|
||||
return None
|
||||
return DifyShellCliToolConfig(name=name if isinstance(name, str) else None, install_commands=commands)
|
||||
|
||||
|
||||
def _shell_env_var(item: object) -> DifyShellEnvVarConfig | None:
|
||||
data = _plain_mapping(item)
|
||||
name = _name_from_mapping(data)
|
||||
if name is None:
|
||||
return None
|
||||
value = data.get("value", data.get("default", ""))
|
||||
if not isinstance(value, str):
|
||||
value = str(value)
|
||||
return DifyShellEnvVarConfig(name=name, value=value)
|
||||
|
||||
|
||||
def _shell_secret_ref(item: object) -> DifyShellSecretRefConfig | None:
|
||||
data = _plain_mapping(item)
|
||||
name = _name_from_mapping(data)
|
||||
if name is None:
|
||||
return None
|
||||
ref = data.get("ref") or data.get("id") or data.get("credential_id") or data.get("provider_credential_id")
|
||||
return DifyShellSecretRefConfig(name=name, ref=str(ref) if ref is not None else None)
|
||||
|
||||
|
||||
def _plain_mapping(item: object) -> dict[str, Any]:
|
||||
if isinstance(item, BaseModel):
|
||||
return item.model_dump(mode="python", exclude_none=True, exclude_defaults=True)
|
||||
if isinstance(item, Mapping):
|
||||
return dict(item)
|
||||
return {}
|
||||
|
||||
|
||||
def _name_from_mapping(item: Mapping[str, Any]) -> str | None:
|
||||
for key in ("name", "key", "env_name", "variable"):
|
||||
value = item.get(key)
|
||||
if isinstance(value, str) and value.strip():
|
||||
return value.strip()
|
||||
return None
|
||||
|
||||
@@ -303,8 +303,18 @@ class WorkflowAgentNodeValidator:
|
||||
f"Workflow Agent node {binding.node_id} has duplicate Dify Plugin Tool name {exposed_name}."
|
||||
)
|
||||
exposed_names.add(exposed_name)
|
||||
# CLI tools remain saved-but-not-executed. They are allowed at publish
|
||||
# time so existing Agent Soul drafts are not blocked by a reserved field.
|
||||
|
||||
cli_tool_names: set[str] = set()
|
||||
for cli_tool in agent_soul.tools.cli_tools:
|
||||
name = cli_tool.get("name") or cli_tool.get("tool_name") or cli_tool.get("label")
|
||||
if not isinstance(name, str) or not name.strip():
|
||||
continue
|
||||
normalized_name = name.strip()
|
||||
if normalized_name in cli_tool_names:
|
||||
raise WorkflowAgentNodeValidationError(
|
||||
f"Workflow Agent node {binding.node_id} has duplicate CLI Tool name {normalized_name}."
|
||||
)
|
||||
cli_tool_names.add(normalized_name)
|
||||
|
||||
@staticmethod
|
||||
def _validate_file_ref(
|
||||
|
||||
@@ -1103,6 +1103,70 @@ List workflow apps that reference this Agent App's bound Agent (read-only)
|
||||
| 200 | Referencing workflows listed successfully | [AgentReferencingWorkflowsResponse](#agentreferencingworkflowsresponse) |
|
||||
| 404 | App not found | |
|
||||
|
||||
### /apps/{app_id}/agent-workspace/files
|
||||
|
||||
#### GET
|
||||
##### Description
|
||||
|
||||
List a directory in an Agent App conversation's sandbox workspace (read-only)
|
||||
|
||||
##### Parameters
|
||||
|
||||
| Name | Located in | Description | Required | Schema |
|
||||
| ---- | ---------- | ----------- | -------- | ------ |
|
||||
| app_id | path | Application ID | Yes | string |
|
||||
| conversation_id | query | Agent App conversation ID | Yes | string |
|
||||
| path | query | Directory path relative to the sandbox workspace | No | string |
|
||||
|
||||
##### Responses
|
||||
|
||||
| Code | Description | Schema |
|
||||
| ---- | ----------- | ------ |
|
||||
| 200 | Listing returned | [WorkspaceListResponse](#workspacelistresponse) |
|
||||
|
||||
### /apps/{app_id}/agent-workspace/files/download
|
||||
|
||||
#### GET
|
||||
##### Description
|
||||
|
||||
Download a file from an Agent App conversation's sandbox workspace (read-only)
|
||||
|
||||
##### Parameters
|
||||
|
||||
| Name | Located in | Description | Required | Schema |
|
||||
| ---- | ---------- | ----------- | -------- | ------ |
|
||||
| app_id | path | Application ID | Yes | string |
|
||||
| conversation_id | query | Agent App conversation ID | Yes | string |
|
||||
| path | query | File path relative to the sandbox workspace | Yes | string |
|
||||
|
||||
##### Responses
|
||||
|
||||
| Code | Description | Schema |
|
||||
| ---- | ----------- | ------ |
|
||||
| 200 | File bytes | binary |
|
||||
| 413 | File exceeds the workspace download limit | |
|
||||
|
||||
### /apps/{app_id}/agent-workspace/files/preview
|
||||
|
||||
#### GET
|
||||
##### Description
|
||||
|
||||
Preview a text/binary file in an Agent App conversation's sandbox workspace
|
||||
|
||||
##### Parameters
|
||||
|
||||
| Name | Located in | Description | Required | Schema |
|
||||
| ---- | ---------- | ----------- | -------- | ------ |
|
||||
| app_id | path | Application ID | Yes | string |
|
||||
| conversation_id | query | Agent App conversation ID | Yes | string |
|
||||
| path | query | File path relative to the sandbox workspace | Yes | string |
|
||||
|
||||
##### Responses
|
||||
|
||||
| Code | Description | Schema |
|
||||
| ---- | ----------- | ------ |
|
||||
| 200 | Preview returned | [WorkspacePreviewResponse](#workspacepreviewresponse) |
|
||||
|
||||
### /apps/{app_id}/agent/logs
|
||||
|
||||
#### GET
|
||||
@@ -2623,6 +2687,76 @@ Get workflow run node execution list
|
||||
| 200 | Node executions retrieved successfully | [WorkflowRunNodeExecutionListResponse](#workflowrunnodeexecutionlistresponse) |
|
||||
| 404 | Workflow run not found | |
|
||||
|
||||
### /apps/{app_id}/workflow-runs/{workflow_run_id}/agent-nodes/{node_id}/workspace/files
|
||||
|
||||
#### GET
|
||||
##### Description
|
||||
|
||||
List a directory in a Workflow Agent node's sandbox workspace (read-only)
|
||||
|
||||
##### Parameters
|
||||
|
||||
| Name | Located in | Description | Required | Schema |
|
||||
| ---- | ---------- | ----------- | -------- | ------ |
|
||||
| app_id | path | Application ID | Yes | string |
|
||||
| node_id | path | Workflow Agent node ID | Yes | string |
|
||||
| workflow_run_id | path | Workflow run ID | Yes | string |
|
||||
| node_execution_id | query | Optional workflow node execution ID. When omitted, the latest active session for the node is used. | No | string |
|
||||
| path | query | Directory path relative to the sandbox workspace | No | string |
|
||||
|
||||
##### Responses
|
||||
|
||||
| Code | Description | Schema |
|
||||
| ---- | ----------- | ------ |
|
||||
| 200 | Listing returned | [WorkspaceListResponse](#workspacelistresponse) |
|
||||
|
||||
### /apps/{app_id}/workflow-runs/{workflow_run_id}/agent-nodes/{node_id}/workspace/files/download
|
||||
|
||||
#### GET
|
||||
##### Description
|
||||
|
||||
Download a file from a Workflow Agent node's sandbox workspace (read-only)
|
||||
|
||||
##### Parameters
|
||||
|
||||
| Name | Located in | Description | Required | Schema |
|
||||
| ---- | ---------- | ----------- | -------- | ------ |
|
||||
| app_id | path | Application ID | Yes | string |
|
||||
| node_id | path | Workflow Agent node ID | Yes | string |
|
||||
| workflow_run_id | path | Workflow run ID | Yes | string |
|
||||
| node_execution_id | query | Optional workflow node execution ID. When omitted, the latest active session for the node is used. | No | string |
|
||||
| path | query | File path relative to the sandbox workspace | Yes | string |
|
||||
|
||||
##### Responses
|
||||
|
||||
| Code | Description | Schema |
|
||||
| ---- | ----------- | ------ |
|
||||
| 200 | File bytes | binary |
|
||||
| 413 | File exceeds the workspace download limit | |
|
||||
|
||||
### /apps/{app_id}/workflow-runs/{workflow_run_id}/agent-nodes/{node_id}/workspace/files/preview
|
||||
|
||||
#### GET
|
||||
##### Description
|
||||
|
||||
Preview a text/binary file in a Workflow Agent node's sandbox workspace
|
||||
|
||||
##### Parameters
|
||||
|
||||
| Name | Located in | Description | Required | Schema |
|
||||
| ---- | ---------- | ----------- | -------- | ------ |
|
||||
| app_id | path | Application ID | Yes | string |
|
||||
| node_id | path | Workflow Agent node ID | Yes | string |
|
||||
| workflow_run_id | path | Workflow run ID | Yes | string |
|
||||
| node_execution_id | query | Optional workflow node execution ID. When omitted, the latest active session for the node is used. | No | string |
|
||||
| path | query | File path relative to the sandbox workspace | Yes | string |
|
||||
|
||||
##### Responses
|
||||
|
||||
| Code | Description | Schema |
|
||||
| ---- | ----------- | ------ |
|
||||
| 200 | Preview returned | [WorkspacePreviewResponse](#workspacepreviewresponse) |
|
||||
|
||||
### /apps/{app_id}/workflow/comments
|
||||
|
||||
#### GET
|
||||
@@ -16864,6 +16998,15 @@ Workflow tool configuration
|
||||
| remove_webapp_brand | boolean | | No |
|
||||
| replace_webapp_logo | string | | No |
|
||||
|
||||
#### WorkspaceFileEntryResponse
|
||||
|
||||
| Name | Type | Description | Required |
|
||||
| ---- | ---- | ----------- | -------- |
|
||||
| mtime | integer | | Yes |
|
||||
| name | string | | Yes |
|
||||
| size | integer | | Yes |
|
||||
| type | string | *Enum:* `"dir"`, `"file"`, `"symlink"` | Yes |
|
||||
|
||||
#### WorkspaceInfoPayload
|
||||
|
||||
| Name | Type | Description | Required |
|
||||
@@ -16877,6 +17020,14 @@ Workflow tool configuration
|
||||
| limit | integer | | No |
|
||||
| page | integer | | No |
|
||||
|
||||
#### WorkspaceListResponse
|
||||
|
||||
| Name | Type | Description | Required |
|
||||
| ---- | ---- | ----------- | -------- |
|
||||
| entries | [ [WorkspaceFileEntryResponse](#workspacefileentryresponse) ] | | No |
|
||||
| path | string | | Yes |
|
||||
| truncated | boolean | | No |
|
||||
|
||||
#### WorkspacePermissionResponse
|
||||
|
||||
| Name | Type | Description | Required |
|
||||
@@ -16885,6 +17036,16 @@ Workflow tool configuration
|
||||
| allow_owner_transfer | boolean | | Yes |
|
||||
| workspace_id | string | | Yes |
|
||||
|
||||
#### WorkspacePreviewResponse
|
||||
|
||||
| Name | Type | Description | Required |
|
||||
| ---- | ---- | ----------- | -------- |
|
||||
| binary | boolean | | Yes |
|
||||
| path | string | | Yes |
|
||||
| size | integer | | Yes |
|
||||
| text | string | | No |
|
||||
| truncated | boolean | | Yes |
|
||||
|
||||
#### _AnonymousInlineModel_b1954337d565
|
||||
|
||||
| Name | Type | Description | Required |
|
||||
|
||||
@@ -0,0 +1,220 @@
|
||||
"""Resolve and proxy read-only access to an Agent App conversation's sandbox.
|
||||
|
||||
The Agent App's shell layer runs bash in a per-conversation sandbox workspace on
|
||||
the agent backend. The workspace identity (``session_id``) is generated inside
|
||||
the shell layer and rides the conversation's ``session_snapshot``. This service
|
||||
extracts that id and proxies list/preview/download to the agent backend's
|
||||
read-only workspace endpoints, so the console can show a "sandbox file system"
|
||||
inspector without the API ever touching shellctl directly.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from collections.abc import Callable
|
||||
|
||||
from agenton.compositor import CompositorSessionSnapshot
|
||||
from sqlalchemy import select
|
||||
|
||||
from clients.agent_backend.request_builder import DIFY_SHELL_LAYER_ID
|
||||
from clients.agent_backend.workspace_files_client import (
|
||||
WorkspaceDownloadResult,
|
||||
WorkspaceFilesBackendClient,
|
||||
WorkspaceListResult,
|
||||
WorkspacePreviewResult,
|
||||
)
|
||||
from configs import dify_config
|
||||
from core.app.apps.agent_app.session_store import AgentAppRuntimeSessionStore
|
||||
from core.db.session_factory import session_factory
|
||||
from models.agent import (
|
||||
AgentRuntimeSessionOwnerType,
|
||||
WorkflowAgentRuntimeSession,
|
||||
WorkflowAgentRuntimeSessionStatus,
|
||||
)
|
||||
|
||||
|
||||
class AgentWorkspaceInspectorError(Exception):
|
||||
"""A workspace inspection failure mapped to an HTTP status by the controller."""
|
||||
|
||||
code: str
|
||||
message: str
|
||||
status_code: int
|
||||
|
||||
def __init__(self, code: str, message: str, *, status_code: int = 400) -> None:
|
||||
super().__init__(message)
|
||||
self.code = code
|
||||
self.message = message
|
||||
self.status_code = status_code
|
||||
|
||||
|
||||
class AgentAppWorkspaceService:
|
||||
"""List/preview/download files in an Agent App conversation's sandbox workspace."""
|
||||
|
||||
def __init__(
|
||||
self,
|
||||
*,
|
||||
session_store: AgentAppRuntimeSessionStore | None = None,
|
||||
client_factory: Callable[[], WorkspaceFilesBackendClient] | None = None,
|
||||
) -> None:
|
||||
self._session_store = session_store or AgentAppRuntimeSessionStore()
|
||||
self._client_factory = client_factory or _default_client_factory
|
||||
|
||||
def list_files(self, *, tenant_id: str, app_id: str, conversation_id: str, path: str) -> WorkspaceListResult:
|
||||
session_id = self._resolve_session_id(tenant_id=tenant_id, app_id=app_id, conversation_id=conversation_id)
|
||||
return self._client_factory().list_files(session_id, path)
|
||||
|
||||
def preview(self, *, tenant_id: str, app_id: str, conversation_id: str, path: str) -> WorkspacePreviewResult:
|
||||
session_id = self._resolve_session_id(tenant_id=tenant_id, app_id=app_id, conversation_id=conversation_id)
|
||||
return self._client_factory().preview(session_id, path)
|
||||
|
||||
def download(self, *, tenant_id: str, app_id: str, conversation_id: str, path: str) -> WorkspaceDownloadResult:
|
||||
session_id = self._resolve_session_id(tenant_id=tenant_id, app_id=app_id, conversation_id=conversation_id)
|
||||
return self._client_factory().download(session_id, path)
|
||||
|
||||
def _resolve_session_id(self, *, tenant_id: str, app_id: str, conversation_id: str) -> str:
|
||||
snapshot = self._session_store.load_active_snapshot_for_conversation(
|
||||
tenant_id=tenant_id, app_id=app_id, conversation_id=conversation_id
|
||||
)
|
||||
if snapshot is None:
|
||||
raise AgentWorkspaceInspectorError(
|
||||
"no_active_session",
|
||||
"this conversation has no active sandbox session yet",
|
||||
status_code=404,
|
||||
)
|
||||
session_id = _shell_session_id(snapshot)
|
||||
if not session_id:
|
||||
raise AgentWorkspaceInspectorError(
|
||||
"no_sandbox",
|
||||
"this conversation's agent has no sandbox workspace",
|
||||
status_code=404,
|
||||
)
|
||||
return session_id
|
||||
|
||||
|
||||
class WorkflowAgentWorkspaceService:
|
||||
"""List/preview/download files in a Workflow Agent node sandbox workspace."""
|
||||
|
||||
def __init__(
|
||||
self,
|
||||
*,
|
||||
client_factory: Callable[[], WorkspaceFilesBackendClient] | None = None,
|
||||
) -> None:
|
||||
self._client_factory = client_factory or _default_client_factory
|
||||
|
||||
def list_files(
|
||||
self,
|
||||
*,
|
||||
tenant_id: str,
|
||||
app_id: str,
|
||||
workflow_run_id: str,
|
||||
node_id: str,
|
||||
node_execution_id: str | None,
|
||||
path: str,
|
||||
) -> WorkspaceListResult:
|
||||
session_id = self._resolve_session_id(
|
||||
tenant_id=tenant_id,
|
||||
app_id=app_id,
|
||||
workflow_run_id=workflow_run_id,
|
||||
node_id=node_id,
|
||||
node_execution_id=node_execution_id,
|
||||
)
|
||||
return self._client_factory().list_files(session_id, path)
|
||||
|
||||
def preview(
|
||||
self,
|
||||
*,
|
||||
tenant_id: str,
|
||||
app_id: str,
|
||||
workflow_run_id: str,
|
||||
node_id: str,
|
||||
node_execution_id: str | None,
|
||||
path: str,
|
||||
) -> WorkspacePreviewResult:
|
||||
session_id = self._resolve_session_id(
|
||||
tenant_id=tenant_id,
|
||||
app_id=app_id,
|
||||
workflow_run_id=workflow_run_id,
|
||||
node_id=node_id,
|
||||
node_execution_id=node_execution_id,
|
||||
)
|
||||
return self._client_factory().preview(session_id, path)
|
||||
|
||||
def download(
|
||||
self,
|
||||
*,
|
||||
tenant_id: str,
|
||||
app_id: str,
|
||||
workflow_run_id: str,
|
||||
node_id: str,
|
||||
node_execution_id: str | None,
|
||||
path: str,
|
||||
) -> WorkspaceDownloadResult:
|
||||
session_id = self._resolve_session_id(
|
||||
tenant_id=tenant_id,
|
||||
app_id=app_id,
|
||||
workflow_run_id=workflow_run_id,
|
||||
node_id=node_id,
|
||||
node_execution_id=node_execution_id,
|
||||
)
|
||||
return self._client_factory().download(session_id, path)
|
||||
|
||||
def _resolve_session_id(
|
||||
self,
|
||||
*,
|
||||
tenant_id: str,
|
||||
app_id: str,
|
||||
workflow_run_id: str,
|
||||
node_id: str,
|
||||
node_execution_id: str | None,
|
||||
) -> str:
|
||||
stmt = select(WorkflowAgentRuntimeSession).where(
|
||||
WorkflowAgentRuntimeSession.owner_type == AgentRuntimeSessionOwnerType.WORKFLOW_RUN,
|
||||
WorkflowAgentRuntimeSession.tenant_id == tenant_id,
|
||||
WorkflowAgentRuntimeSession.app_id == app_id,
|
||||
WorkflowAgentRuntimeSession.workflow_run_id == workflow_run_id,
|
||||
WorkflowAgentRuntimeSession.node_id == node_id,
|
||||
WorkflowAgentRuntimeSession.status == WorkflowAgentRuntimeSessionStatus.ACTIVE,
|
||||
)
|
||||
if node_execution_id:
|
||||
stmt = stmt.where(WorkflowAgentRuntimeSession.node_execution_id == node_execution_id)
|
||||
stmt = stmt.order_by(WorkflowAgentRuntimeSession.updated_at.desc()).limit(1)
|
||||
|
||||
with session_factory.create_session() as session:
|
||||
row = session.scalar(stmt)
|
||||
|
||||
if row is None:
|
||||
raise AgentWorkspaceInspectorError(
|
||||
"no_active_session",
|
||||
"this workflow Agent node has no active sandbox session yet",
|
||||
status_code=404,
|
||||
)
|
||||
snapshot = CompositorSessionSnapshot.model_validate_json(row.session_snapshot)
|
||||
session_id = _shell_session_id(snapshot)
|
||||
if not session_id:
|
||||
raise AgentWorkspaceInspectorError(
|
||||
"no_sandbox",
|
||||
"this workflow Agent node has no sandbox workspace",
|
||||
status_code=404,
|
||||
)
|
||||
return session_id
|
||||
|
||||
|
||||
def _shell_session_id(snapshot: CompositorSessionSnapshot) -> str | None:
|
||||
for layer in snapshot.layers:
|
||||
if layer.name == DIFY_SHELL_LAYER_ID:
|
||||
session_id = layer.runtime_state.get("session_id")
|
||||
return session_id if isinstance(session_id, str) and session_id else None
|
||||
return None
|
||||
|
||||
|
||||
def _default_client_factory() -> WorkspaceFilesBackendClient:
|
||||
base_url = dify_config.AGENT_BACKEND_BASE_URL
|
||||
if not base_url:
|
||||
raise AgentWorkspaceInspectorError(
|
||||
"inspector_unavailable",
|
||||
"the sandbox file inspector is not available (agent backend not configured)",
|
||||
status_code=503,
|
||||
)
|
||||
return WorkspaceFilesBackendClient(base_url)
|
||||
|
||||
|
||||
__all__ = ["AgentAppWorkspaceService", "AgentWorkspaceInspectorError", "WorkflowAgentWorkspaceService"]
|
||||
@@ -16,6 +16,7 @@ from dify_agent.layers.dify_plugin import (
|
||||
)
|
||||
from dify_agent.layers.execution_context import DIFY_EXECUTION_CONTEXT_LAYER_TYPE_ID, DifyExecutionContextLayerConfig
|
||||
from dify_agent.layers.output import DIFY_OUTPUT_LAYER_TYPE_ID
|
||||
from dify_agent.layers.shell import DIFY_SHELL_LAYER_TYPE_ID, DifyShellEnvVarConfig, DifyShellLayerConfig
|
||||
from dify_agent.protocol import (
|
||||
DIFY_AGENT_HISTORY_LAYER_ID,
|
||||
DIFY_AGENT_MODEL_LAYER_ID,
|
||||
@@ -30,6 +31,7 @@ from clients.agent_backend import (
|
||||
DIFY_PLUGIN_TOOLS_LAYER_ID,
|
||||
WORKFLOW_NODE_JOB_PROMPT_LAYER_ID,
|
||||
WORKFLOW_USER_PROMPT_LAYER_ID,
|
||||
AgentBackendAgentAppRunInput,
|
||||
AgentBackendModelConfig,
|
||||
AgentBackendOutputConfig,
|
||||
AgentBackendRunRequestBuilder,
|
||||
@@ -37,6 +39,7 @@ from clients.agent_backend import (
|
||||
CleanupLayerSpec,
|
||||
redact_for_agent_backend_log,
|
||||
)
|
||||
from clients.agent_backend.request_builder import DIFY_SHELL_LAYER_ID
|
||||
|
||||
|
||||
def _run_input() -> AgentBackendWorkflowNodeRunInput:
|
||||
@@ -249,3 +252,65 @@ def test_redact_for_agent_backend_log_hides_credentials():
|
||||
redacted = cast(dict[str, Any], redact_for_agent_backend_log(request))
|
||||
|
||||
assert redacted["composition"]["layers"][5]["config"]["credentials"] == "[REDACTED]"
|
||||
|
||||
|
||||
def _agent_app_input(*, include_shell: bool = False) -> AgentBackendAgentAppRunInput:
|
||||
return AgentBackendAgentAppRunInput(
|
||||
model=AgentBackendModelConfig(
|
||||
plugin_id="langgenius/openai",
|
||||
model_provider="openai",
|
||||
model="gpt-test",
|
||||
credentials={"api_key": "secret-key"},
|
||||
),
|
||||
execution_context=DifyExecutionContextLayerConfig(
|
||||
tenant_id="tenant-1",
|
||||
user_id="user-1",
|
||||
conversation_id="conv-1",
|
||||
invoke_from="agent_app",
|
||||
),
|
||||
agent_soul_prompt="You are Iris.",
|
||||
user_prompt="List files.",
|
||||
include_shell=include_shell,
|
||||
metadata={"conversation_id": "conv-1"},
|
||||
)
|
||||
|
||||
|
||||
def test_workflow_request_builder_omits_shell_layer_by_default():
|
||||
request = AgentBackendRunRequestBuilder().build_for_workflow_node(_run_input())
|
||||
assert DIFY_SHELL_LAYER_ID not in {layer.name for layer in request.composition.layers}
|
||||
|
||||
|
||||
def test_workflow_request_builder_adds_shell_layer_when_include_shell():
|
||||
run_input = _run_input()
|
||||
run_input.include_shell = True
|
||||
run_input.shell_config = DifyShellLayerConfig(env=[DifyShellEnvVarConfig(name="PROJECT_NAME", value="demo")])
|
||||
|
||||
request = AgentBackendRunRequestBuilder().build_for_workflow_node(run_input)
|
||||
layers = {layer.name: layer for layer in request.composition.layers}
|
||||
|
||||
assert DIFY_SHELL_LAYER_ID in layers
|
||||
shell = layers[DIFY_SHELL_LAYER_ID]
|
||||
assert shell.type == DIFY_SHELL_LAYER_TYPE_ID
|
||||
# The shell layer declares NoLayerDeps, so the spec must carry no deps.
|
||||
assert not shell.deps
|
||||
shell_config = cast(DifyShellLayerConfig, shell.config)
|
||||
assert shell_config.env[0].name == "PROJECT_NAME"
|
||||
|
||||
|
||||
def test_agent_app_request_builder_omits_shell_layer_by_default():
|
||||
request = AgentBackendRunRequestBuilder().build_for_agent_app(_agent_app_input())
|
||||
assert DIFY_SHELL_LAYER_ID not in {layer.name for layer in request.composition.layers}
|
||||
|
||||
|
||||
def test_agent_app_request_builder_adds_shell_layer_when_include_shell():
|
||||
run_input = _agent_app_input(include_shell=True)
|
||||
run_input.shell_config = DifyShellLayerConfig(env=[DifyShellEnvVarConfig(name="APP_ENV", value="enabled")])
|
||||
|
||||
request = AgentBackendRunRequestBuilder().build_for_agent_app(run_input)
|
||||
layers = {layer.name: layer for layer in request.composition.layers}
|
||||
|
||||
assert DIFY_SHELL_LAYER_ID in layers
|
||||
assert layers[DIFY_SHELL_LAYER_ID].type == DIFY_SHELL_LAYER_TYPE_ID
|
||||
assert not layers[DIFY_SHELL_LAYER_ID].deps
|
||||
shell_config = cast(DifyShellLayerConfig, layers[DIFY_SHELL_LAYER_ID].config)
|
||||
assert shell_config.env[0].name == "APP_ENV"
|
||||
|
||||
@@ -0,0 +1,152 @@
|
||||
"""Unit tests for the API-side workspace files backend client."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import base64
|
||||
import json
|
||||
from collections.abc import Callable
|
||||
|
||||
import httpx
|
||||
import pytest
|
||||
|
||||
from clients.agent_backend.errors import AgentBackendHTTPError, AgentBackendTransportError
|
||||
from clients.agent_backend.workspace_files_client import WorkspaceFilesBackendClient
|
||||
|
||||
|
||||
def _client(handler: Callable[[httpx.Request], httpx.Response]) -> WorkspaceFilesBackendClient:
|
||||
return WorkspaceFilesBackendClient("http://backend", transport=httpx.MockTransport(handler))
|
||||
|
||||
|
||||
def test_list_files_parses_entries() -> None:
|
||||
def handler(request: httpx.Request) -> httpx.Response:
|
||||
assert request.url.path == "/workspaces/abc1234/files"
|
||||
assert request.url.params.get("path") == "sub"
|
||||
return httpx.Response(
|
||||
200,
|
||||
json={
|
||||
"path": "sub",
|
||||
"entries": [{"name": "a.txt", "type": "file", "size": 3, "mtime": 10}],
|
||||
"truncated": False,
|
||||
},
|
||||
)
|
||||
|
||||
result = _client(handler).list_files("abc1234", "sub")
|
||||
|
||||
assert result.path == "sub"
|
||||
assert result.entries[0].name == "a.txt"
|
||||
assert result.entries[0].type == "file"
|
||||
assert result.truncated is False
|
||||
|
||||
|
||||
def test_preview_parses_text() -> None:
|
||||
def handler(request: httpx.Request) -> httpx.Response:
|
||||
assert request.url.path == "/workspaces/abc1234/files/preview"
|
||||
return httpx.Response(
|
||||
200, json={"path": "n.txt", "size": 5, "truncated": False, "binary": False, "text": "hello"}
|
||||
)
|
||||
|
||||
result = _client(handler).preview("abc1234", "n.txt")
|
||||
|
||||
assert result.binary is False
|
||||
assert result.text == "hello"
|
||||
|
||||
|
||||
def test_download_decodes_base64_to_bytes() -> None:
|
||||
raw = bytes(range(64))
|
||||
|
||||
def handler(request: httpx.Request) -> httpx.Response:
|
||||
assert request.url.path == "/workspaces/abc1234/files/download"
|
||||
return httpx.Response(
|
||||
200,
|
||||
json={
|
||||
"path": "b.bin",
|
||||
"size": len(raw),
|
||||
"truncated": False,
|
||||
"content_base64": base64.b64encode(raw).decode(),
|
||||
},
|
||||
)
|
||||
|
||||
result = _client(handler).download("abc1234", "b.bin")
|
||||
|
||||
assert result.content == raw
|
||||
assert result.size == 64
|
||||
|
||||
|
||||
def test_http_error_preserves_status_and_detail() -> None:
|
||||
def handler(request: httpx.Request) -> httpx.Response:
|
||||
return httpx.Response(404, json={"detail": {"code": "not_found", "message": "path not found in workspace"}})
|
||||
|
||||
with pytest.raises(AgentBackendHTTPError) as exc_info:
|
||||
_client(handler).list_files("abc1234", "missing")
|
||||
|
||||
assert exc_info.value.status_code == 404
|
||||
assert exc_info.value.detail == {"code": "not_found", "message": "path not found in workspace"}
|
||||
|
||||
|
||||
def test_http_error_with_non_json_body_uses_response_text() -> None:
|
||||
def handler(request: httpx.Request) -> httpx.Response:
|
||||
return httpx.Response(500, text="backend exploded")
|
||||
|
||||
with pytest.raises(AgentBackendHTTPError) as exc_info:
|
||||
_client(handler).preview("abc1234", "note.txt")
|
||||
|
||||
assert exc_info.value.status_code == 500
|
||||
assert exc_info.value.detail == "backend exploded"
|
||||
|
||||
|
||||
def test_transport_failure_becomes_transport_error() -> None:
|
||||
def handler(request: httpx.Request) -> httpx.Response:
|
||||
raise httpx.ConnectError("connection refused")
|
||||
|
||||
with pytest.raises(AgentBackendTransportError):
|
||||
_client(handler).list_files("abc1234", ".")
|
||||
|
||||
|
||||
def test_download_without_content_is_502() -> None:
|
||||
def handler(request: httpx.Request) -> httpx.Response:
|
||||
return httpx.Response(200, json={"path": "b.bin", "size": 0, "truncated": False})
|
||||
|
||||
with pytest.raises(AgentBackendHTTPError) as exc_info:
|
||||
_client(handler).download("abc1234", "b.bin")
|
||||
|
||||
assert exc_info.value.status_code == 502
|
||||
|
||||
|
||||
def test_download_with_invalid_base64_is_502() -> None:
|
||||
def handler(request: httpx.Request) -> httpx.Response:
|
||||
return httpx.Response(200, json={"path": "b.bin", "size": 3, "truncated": False, "content_base64": "not-@@@"})
|
||||
|
||||
with pytest.raises(AgentBackendHTTPError) as exc_info:
|
||||
_client(handler).download("abc1234", "b.bin")
|
||||
|
||||
assert exc_info.value.status_code == 502
|
||||
|
||||
|
||||
def test_download_uses_decoded_size_when_backend_size_is_invalid() -> None:
|
||||
raw = b"abc"
|
||||
|
||||
def handler(request: httpx.Request) -> httpx.Response:
|
||||
return httpx.Response(
|
||||
200,
|
||||
json={
|
||||
"path": "b.bin",
|
||||
"size": "unknown",
|
||||
"truncated": True,
|
||||
"content_base64": base64.b64encode(raw).decode(),
|
||||
},
|
||||
)
|
||||
|
||||
result = _client(handler).download("abc1234", "b.bin")
|
||||
|
||||
assert result.size == len(raw)
|
||||
assert result.truncated is True
|
||||
|
||||
|
||||
def test_non_object_body_is_502() -> None:
|
||||
def handler(request: httpx.Request) -> httpx.Response:
|
||||
return httpx.Response(200, content=json.dumps([1, 2, 3]), headers={"content-type": "application/json"})
|
||||
|
||||
with pytest.raises(AgentBackendHTTPError) as exc_info:
|
||||
_client(handler).list_files("abc1234", ".")
|
||||
|
||||
assert exc_info.value.status_code == 502
|
||||
@@ -0,0 +1,199 @@
|
||||
from __future__ import annotations
|
||||
|
||||
from types import SimpleNamespace
|
||||
|
||||
import pytest
|
||||
|
||||
from clients.agent_backend.errors import AgentBackendHTTPError, AgentBackendTransportError
|
||||
from clients.agent_backend.workspace_files_client import (
|
||||
WorkspaceDownloadResult,
|
||||
WorkspaceFileEntry,
|
||||
WorkspaceListResult,
|
||||
WorkspacePreviewResult,
|
||||
)
|
||||
from controllers.console import agent_app_workspace as module
|
||||
from services.agent_app_workspace_service import AgentWorkspaceInspectorError
|
||||
|
||||
|
||||
def _unwrapped_get(resource_cls):
|
||||
func = resource_cls.get
|
||||
while hasattr(func, "__wrapped__"):
|
||||
func = func.__wrapped__
|
||||
return func
|
||||
|
||||
|
||||
class _AgentAppService:
|
||||
def __init__(self) -> None:
|
||||
self.calls: list[tuple[str, str, str, str, str]] = []
|
||||
|
||||
def list_files(self, *, tenant_id: str, app_id: str, conversation_id: str, path: str) -> WorkspaceListResult:
|
||||
self.calls.append(("list", tenant_id, app_id, conversation_id, path))
|
||||
return WorkspaceListResult(
|
||||
path=path,
|
||||
entries=[WorkspaceFileEntry(name="a.txt", type="file", size=3, mtime=10)],
|
||||
truncated=False,
|
||||
)
|
||||
|
||||
def preview(self, *, tenant_id: str, app_id: str, conversation_id: str, path: str) -> WorkspacePreviewResult:
|
||||
self.calls.append(("preview", tenant_id, app_id, conversation_id, path))
|
||||
return WorkspacePreviewResult(path=path, size=5, truncated=False, binary=False, text="hello")
|
||||
|
||||
def download(self, *, tenant_id: str, app_id: str, conversation_id: str, path: str) -> WorkspaceDownloadResult:
|
||||
self.calls.append(("download", tenant_id, app_id, conversation_id, path))
|
||||
return WorkspaceDownloadResult(path=path, size=3, truncated=False, content=b"abc")
|
||||
|
||||
|
||||
class _WorkflowService:
|
||||
def __init__(self) -> None:
|
||||
self.calls: list[tuple[str, str, str, str, str, str | None, str]] = []
|
||||
|
||||
def list_files(
|
||||
self,
|
||||
*,
|
||||
tenant_id: str,
|
||||
app_id: str,
|
||||
workflow_run_id: str,
|
||||
node_id: str,
|
||||
node_execution_id: str | None,
|
||||
path: str,
|
||||
) -> WorkspaceListResult:
|
||||
self.calls.append(("list", tenant_id, app_id, workflow_run_id, node_id, node_execution_id, path))
|
||||
return WorkspaceListResult(path=path, entries=[], truncated=False)
|
||||
|
||||
def preview(
|
||||
self,
|
||||
*,
|
||||
tenant_id: str,
|
||||
app_id: str,
|
||||
workflow_run_id: str,
|
||||
node_id: str,
|
||||
node_execution_id: str | None,
|
||||
path: str,
|
||||
) -> WorkspacePreviewResult:
|
||||
self.calls.append(("preview", tenant_id, app_id, workflow_run_id, node_id, node_execution_id, path))
|
||||
return WorkspacePreviewResult(path=path, size=5, truncated=False, binary=False, text="hello")
|
||||
|
||||
def download(
|
||||
self,
|
||||
*,
|
||||
tenant_id: str,
|
||||
app_id: str,
|
||||
workflow_run_id: str,
|
||||
node_id: str,
|
||||
node_execution_id: str | None,
|
||||
path: str,
|
||||
) -> WorkspaceDownloadResult:
|
||||
self.calls.append(("download", tenant_id, app_id, workflow_run_id, node_id, node_execution_id, path))
|
||||
return WorkspaceDownloadResult(path=path, size=3, truncated=False, content=b"abc")
|
||||
|
||||
|
||||
def test_handle_maps_workspace_and_agent_backend_errors() -> None:
|
||||
assert module._handle(AgentWorkspaceInspectorError("no_sandbox", "no sandbox", status_code=404)) == (
|
||||
{"code": "no_sandbox", "message": "no sandbox"},
|
||||
404,
|
||||
)
|
||||
assert module._handle(
|
||||
AgentBackendHTTPError("not found", status_code=404, detail={"code": "not_found", "message": "missing"})
|
||||
) == ({"code": "not_found", "message": "missing"}, 404)
|
||||
assert module._handle(AgentBackendHTTPError("bad", status_code=500, detail="backend exploded")) == (
|
||||
{"code": "agent_backend_error", "message": "backend exploded"},
|
||||
500,
|
||||
)
|
||||
assert module._handle(AgentBackendTransportError("connection refused")) == (
|
||||
{"code": "agent_backend_unreachable", "message": "connection refused"},
|
||||
502,
|
||||
)
|
||||
with pytest.raises(RuntimeError):
|
||||
module._handle(RuntimeError("boom"))
|
||||
|
||||
|
||||
def test_download_response_returns_binary_or_too_large_error() -> None:
|
||||
response = module._download_response(
|
||||
WorkspaceDownloadResult(path="dir/report.txt", size=3, truncated=False, content=b"abc")
|
||||
)
|
||||
|
||||
assert response.status_code == 200
|
||||
assert response.data == b"abc"
|
||||
assert response.headers["Content-Disposition"] == 'attachment; filename="report.txt"'
|
||||
assert response.headers["Content-Length"] == "3"
|
||||
assert response.headers["X-Workspace-File-Size"] == "3"
|
||||
|
||||
assert module._download_response(WorkspaceDownloadResult(path="", size=10, truncated=True, content=b"")) == (
|
||||
{
|
||||
"code": "workspace_file_too_large",
|
||||
"message": (
|
||||
"file exceeds the workspace download limit; use preview for partial text or download a smaller file"
|
||||
),
|
||||
"size": 10,
|
||||
},
|
||||
413,
|
||||
)
|
||||
|
||||
|
||||
def test_agent_app_workspace_resources_proxy_service(monkeypatch: pytest.MonkeyPatch) -> None:
|
||||
service = _AgentAppService()
|
||||
monkeypatch.setattr(module, "AgentAppWorkspaceService", lambda: service)
|
||||
monkeypatch.setattr(module, "current_account_with_tenant", lambda: (None, "tenant-1"))
|
||||
monkeypatch.setattr(
|
||||
module,
|
||||
"query_params_from_request",
|
||||
lambda model: SimpleNamespace(conversation_id="conv-1", path="sub/report.txt"),
|
||||
)
|
||||
app_model = SimpleNamespace(id="app-1")
|
||||
|
||||
listing = _unwrapped_get(module.AgentAppWorkspaceListResource)(object(), app_model)
|
||||
preview = _unwrapped_get(module.AgentAppWorkspacePreviewResource)(object(), app_model)
|
||||
download = _unwrapped_get(module.AgentAppWorkspaceDownloadResource)(object(), app_model)
|
||||
|
||||
assert listing["entries"][0]["name"] == "a.txt"
|
||||
assert preview["text"] == "hello"
|
||||
assert download.data == b"abc"
|
||||
assert service.calls == [
|
||||
("list", "tenant-1", "app-1", "conv-1", "sub/report.txt"),
|
||||
("preview", "tenant-1", "app-1", "conv-1", "sub/report.txt"),
|
||||
("download", "tenant-1", "app-1", "conv-1", "sub/report.txt"),
|
||||
]
|
||||
|
||||
|
||||
def test_agent_app_workspace_resource_returns_normalized_errors(monkeypatch: pytest.MonkeyPatch) -> None:
|
||||
class FailingService:
|
||||
def list_files(self, **kwargs):
|
||||
raise AgentWorkspaceInspectorError("no_active_session", "no active session", status_code=404)
|
||||
|
||||
monkeypatch.setattr(module, "AgentAppWorkspaceService", FailingService)
|
||||
monkeypatch.setattr(module, "current_account_with_tenant", lambda: (None, "tenant-1"))
|
||||
monkeypatch.setattr(
|
||||
module,
|
||||
"query_params_from_request",
|
||||
lambda model: SimpleNamespace(conversation_id="conv-1", path="."),
|
||||
)
|
||||
|
||||
assert _unwrapped_get(module.AgentAppWorkspaceListResource)(object(), SimpleNamespace(id="app-1")) == (
|
||||
{"code": "no_active_session", "message": "no active session"},
|
||||
404,
|
||||
)
|
||||
|
||||
|
||||
def test_workflow_agent_workspace_resources_proxy_service(monkeypatch: pytest.MonkeyPatch) -> None:
|
||||
service = _WorkflowService()
|
||||
monkeypatch.setattr(module, "WorkflowAgentWorkspaceService", lambda: service)
|
||||
monkeypatch.setattr(module, "current_account_with_tenant", lambda: (None, "tenant-1"))
|
||||
monkeypatch.setattr(
|
||||
module,
|
||||
"query_params_from_request",
|
||||
lambda model: SimpleNamespace(node_execution_id="exec-1", path="out.txt"),
|
||||
)
|
||||
app_model = SimpleNamespace(id="app-1")
|
||||
|
||||
listing = _unwrapped_get(module.WorkflowAgentWorkspaceListResource)(object(), app_model, "run-1", "agent-node")
|
||||
preview = _unwrapped_get(module.WorkflowAgentWorkspacePreviewResource)(object(), app_model, "run-1", "agent-node")
|
||||
download = _unwrapped_get(module.WorkflowAgentWorkspaceDownloadResource)(object(), app_model, "run-1", "agent-node")
|
||||
|
||||
assert listing["path"] == "out.txt"
|
||||
assert preview["text"] == "hello"
|
||||
assert download.data == b"abc"
|
||||
assert service.calls == [
|
||||
("list", "tenant-1", "app-1", "run-1", "agent-node", "exec-1", "out.txt"),
|
||||
("preview", "tenant-1", "app-1", "run-1", "agent-node", "exec-1", "out.txt"),
|
||||
("download", "tenant-1", "app-1", "run-1", "agent-node", "exec-1", "out.txt"),
|
||||
]
|
||||
@@ -14,6 +14,7 @@ from clients.agent_backend import (
|
||||
AgentBackendModelConfig,
|
||||
AgentBackendRunRequestBuilder,
|
||||
)
|
||||
from clients.agent_backend.request_builder import DIFY_SHELL_LAYER_ID
|
||||
from core.app.apps.agent_app.runtime_request_builder import (
|
||||
AgentAppRuntimeBuildContext,
|
||||
AgentAppRuntimeRequestBuilder,
|
||||
@@ -142,3 +143,37 @@ class TestAgentAppRuntimeRequestBuilder:
|
||||
with pytest.raises(AgentAppRuntimeRequestBuildError) as exc:
|
||||
builder.build(_ctx(AgentSoulConfig()))
|
||||
assert exc.value.error_code == "agent_model_not_configured"
|
||||
|
||||
def test_build_maps_agent_soul_shell_settings_to_shell_layer(self, monkeypatch: pytest.MonkeyPatch):
|
||||
monkeypatch.setattr("core.app.apps.agent_app.runtime_request_builder.dify_config.AGENT_SHELL_ENABLED", True)
|
||||
soul = AgentSoulConfig.model_validate(
|
||||
{
|
||||
"model": {
|
||||
"plugin_id": "langgenius/openai",
|
||||
"model_provider": "langgenius/openai/openai",
|
||||
"model": "gpt-4o-mini",
|
||||
},
|
||||
"tools": {"cli_tools": [{"name": "ripgrep", "install_command": "apt-get install -y ripgrep"}]},
|
||||
"env": {"variables": [{"name": "PROJECT_NAME", "value": "demo"}]},
|
||||
"sandbox": {"provider": "independent", "config": {"cpu": 2}},
|
||||
}
|
||||
)
|
||||
builder = AgentAppRuntimeRequestBuilder(
|
||||
credentials_provider=_FakeCredentialsProvider(),
|
||||
plugin_tools_builder=_NoToolsBuilder(), # type: ignore[arg-type]
|
||||
)
|
||||
|
||||
result = builder.build(_ctx(soul))
|
||||
|
||||
dumped = result.request.model_dump(mode="json")
|
||||
shell_config = {layer["name"]: layer for layer in dumped["composition"]["layers"]}[DIFY_SHELL_LAYER_ID][
|
||||
"config"
|
||||
]
|
||||
assert shell_config["cli_tools"][0]["install_commands"] == ["apt-get install -y ripgrep"]
|
||||
assert shell_config["env"][0] == {"name": "PROJECT_NAME", "value": "demo"}
|
||||
assert shell_config["sandbox"] == {"provider": "independent", "config": {"cpu": 2}}
|
||||
assert result.metadata["agent_tools"] == {
|
||||
"dify_tool_count": 0,
|
||||
"dify_tool_names": [],
|
||||
"cli_tool_count": 1,
|
||||
}
|
||||
|
||||
@@ -119,7 +119,7 @@ def test_distinct_conversations_do_not_collide():
|
||||
assert session.query(AgentRuntimeSession).count() == 2
|
||||
|
||||
|
||||
def test_distinct_agent_config_snapshots_do_not_reuse_prior_session():
|
||||
def test_distinct_agent_config_snapshots_keep_only_latest_active_session():
|
||||
store = AgentAppRuntimeSessionStore()
|
||||
store.save_active_snapshot(
|
||||
scope=_scope(agent_config_snapshot_id="snap-1"),
|
||||
@@ -130,9 +130,71 @@ def test_distinct_agent_config_snapshots_do_not_reuse_prior_session():
|
||||
scope=_scope(agent_config_snapshot_id="snap-2"), backend_run_id="b", snapshot=_snapshot(messages=2)
|
||||
)
|
||||
|
||||
assert store.load_active_snapshot(_scope(agent_config_snapshot_id="snap-1")) is not None
|
||||
assert store.load_active_snapshot(_scope(agent_config_snapshot_id="snap-1")) is None
|
||||
assert store.load_active_snapshot(_scope(agent_config_snapshot_id="snap-2")) is not None
|
||||
with session_factory.create_session() as session:
|
||||
rows = session.query(AgentRuntimeSession).order_by(AgentRuntimeSession.backend_run_id).all()
|
||||
assert len(rows) == 2
|
||||
assert [row.agent_config_snapshot_id for row in rows] == ["snap-1", "snap-2"]
|
||||
assert [row.status for row in rows] == [AgentRuntimeSessionStatus.CLEANED, AgentRuntimeSessionStatus.ACTIVE]
|
||||
|
||||
|
||||
def test_load_for_conversation_resolves_without_agent_or_config_scope():
|
||||
store = AgentAppRuntimeSessionStore()
|
||||
store.save_active_snapshot(scope=_scope(), backend_run_id="run-1", snapshot=_snapshot(messages=2))
|
||||
|
||||
# The inspector only knows tenant/app/conversation, not the agent config version.
|
||||
loaded = store.load_active_snapshot_for_conversation(tenant_id="tenant-1", app_id="app-1", conversation_id="conv-1")
|
||||
assert loaded is not None
|
||||
assert loaded.layers[0].runtime_state["messages"] == [
|
||||
{"role": "user", "content": "m0"},
|
||||
{"role": "user", "content": "m1"},
|
||||
]
|
||||
|
||||
|
||||
def test_load_for_conversation_uses_latest_active_snapshot_after_config_change():
|
||||
store = AgentAppRuntimeSessionStore()
|
||||
store.save_active_snapshot(
|
||||
scope=_scope(agent_config_snapshot_id="snap-1"), backend_run_id="a", snapshot=_snapshot()
|
||||
)
|
||||
store.save_active_snapshot(
|
||||
scope=_scope(agent_config_snapshot_id="snap-2"), backend_run_id="b", snapshot=_snapshot(messages=3)
|
||||
)
|
||||
|
||||
loaded = store.load_active_snapshot_for_conversation(tenant_id="tenant-1", app_id="app-1", conversation_id="conv-1")
|
||||
|
||||
assert loaded is not None
|
||||
assert loaded.layers[0].runtime_state["messages"] == [
|
||||
{"role": "user", "content": "m0"},
|
||||
{"role": "user", "content": "m1"},
|
||||
{"role": "user", "content": "m2"},
|
||||
]
|
||||
|
||||
|
||||
def test_load_for_conversation_returns_none_when_cleaned_or_absent():
|
||||
store = AgentAppRuntimeSessionStore()
|
||||
assert (
|
||||
store.load_active_snapshot_for_conversation(tenant_id="tenant-1", app_id="app-1", conversation_id="conv-1")
|
||||
is None
|
||||
)
|
||||
|
||||
store.save_active_snapshot(scope=_scope(), backend_run_id="run-1", snapshot=_snapshot())
|
||||
store.mark_cleaned(scope=_scope(), backend_run_id="cleanup-1")
|
||||
assert (
|
||||
store.load_active_snapshot_for_conversation(tenant_id="tenant-1", app_id="app-1", conversation_id="conv-1")
|
||||
is None
|
||||
)
|
||||
|
||||
|
||||
def test_load_for_conversation_isolates_other_conversations():
|
||||
store = AgentAppRuntimeSessionStore()
|
||||
store.save_active_snapshot(scope=_scope(conversation_id="conv-A"), backend_run_id="a", snapshot=_snapshot())
|
||||
|
||||
assert (
|
||||
store.load_active_snapshot_for_conversation(tenant_id="tenant-1", app_id="app-1", conversation_id="conv-B")
|
||||
is None
|
||||
)
|
||||
assert (
|
||||
store.load_active_snapshot_for_conversation(tenant_id="tenant-1", app_id="app-1", conversation_id="conv-A")
|
||||
is not None
|
||||
)
|
||||
|
||||
@@ -7,12 +7,14 @@ from dify_agent.layers.dify_plugin import DifyPluginToolConfig, DifyPluginToolsL
|
||||
from dify_agent.protocol import DIFY_AGENT_HISTORY_LAYER_ID, DIFY_AGENT_MODEL_LAYER_ID
|
||||
|
||||
from clients.agent_backend import DIFY_EXECUTION_CONTEXT_LAYER_ID, DIFY_PLUGIN_TOOLS_LAYER_ID
|
||||
from clients.agent_backend.request_builder import DIFY_SHELL_LAYER_ID
|
||||
from core.app.entities.app_invoke_entities import DifyRunContext, InvokeFrom, UserFrom
|
||||
from core.workflow.nodes.agent_v2.plugin_tools_builder import WorkflowAgentPluginToolsBuilder
|
||||
from core.workflow.nodes.agent_v2.runtime_request_builder import (
|
||||
WorkflowAgentRuntimeBuildContext,
|
||||
WorkflowAgentRuntimeRequestBuilder,
|
||||
WorkflowAgentRuntimeRequestBuildError,
|
||||
build_shell_layer_config,
|
||||
)
|
||||
from graphon.variables.segments import StringSegment
|
||||
from models.agent import Agent, AgentConfigSnapshot, WorkflowAgentNodeBinding
|
||||
@@ -224,13 +226,93 @@ def test_builds_workflow_run_request_with_file_output_schema_and_reserved_metada
|
||||
assert dumped["idempotency_key"] == "node-exec-1"
|
||||
output_schema = dumped["composition"]["layers"][-1]["config"]["json_schema"]
|
||||
assert output_schema["properties"]["report"]["properties"]["file_id"]["type"] == "string"
|
||||
assert output_schema["properties"]["report"]["required"] == ["file_id"]
|
||||
assert output_schema["properties"]["confidence"]["type"] == "number"
|
||||
assert output_schema["required"] == ["report"]
|
||||
assert dumped["composition"]["layers"][5]["config"]["model_settings"] == {"temperature": 0.2}
|
||||
assert result.metadata["runtime_support"]["reserved_status"]["tools.dify_tools"] == "supported_when_config_valid"
|
||||
assert result.metadata["runtime_support"]["reserved_status"]["tools.cli_tools"] == "reserved_not_executed"
|
||||
warnings = result.metadata["runtime_support"]["unsupported_runtime_warnings"]
|
||||
assert warnings[0]["section"] == "agent_soul.tools.cli_tools"
|
||||
assert result.metadata["runtime_support"]["reserved_status"]["tools.cli_tools"] == "supported_by_shell_bootstrap"
|
||||
assert result.metadata["runtime_support"]["unsupported_runtime_warnings"] == []
|
||||
|
||||
|
||||
def test_build_maps_agent_soul_shell_settings_to_shell_layer(monkeypatch: pytest.MonkeyPatch):
|
||||
monkeypatch.setattr("core.workflow.nodes.agent_v2.runtime_request_builder.dify_config.AGENT_SHELL_ENABLED", True)
|
||||
context = _context()
|
||||
snapshot = AgentConfigSnapshot(
|
||||
id="snapshot-1",
|
||||
tenant_id="tenant-1",
|
||||
agent_id="agent-1",
|
||||
version=1,
|
||||
config_snapshot=AgentSoulConfig(
|
||||
prompt={"system_prompt": "You are careful."},
|
||||
model=AgentSoulModelConfig(
|
||||
plugin_id="langgenius/openai",
|
||||
model_provider="openai",
|
||||
model="gpt-test",
|
||||
),
|
||||
tools={"cli_tools": [{"name": "ripgrep", "install_commands": ["apt-get install -y ripgrep"]}]},
|
||||
env={"variables": [{"name": "PROJECT_NAME", "value": "demo"}]},
|
||||
sandbox={"provider": "independent", "config": {"cpu": 2}},
|
||||
),
|
||||
)
|
||||
context = replace(context, snapshot=snapshot)
|
||||
|
||||
result = WorkflowAgentRuntimeRequestBuilder(credentials_provider=FakeCredentialsProvider()).build(context)
|
||||
|
||||
dumped = result.request.model_dump(mode="json")
|
||||
shell_config = {layer["name"]: layer for layer in dumped["composition"]["layers"]}[DIFY_SHELL_LAYER_ID]["config"]
|
||||
assert shell_config["cli_tools"][0]["install_commands"] == ["apt-get install -y ripgrep"]
|
||||
assert shell_config["env"][0] == {"name": "PROJECT_NAME", "value": "demo"}
|
||||
assert shell_config["sandbox"] == {"provider": "independent", "config": {"cpu": 2}}
|
||||
assert result.metadata["agent_tools"] == {
|
||||
"dify_tool_count": 0,
|
||||
"dify_tool_names": [],
|
||||
"cli_tool_count": 1,
|
||||
}
|
||||
|
||||
|
||||
def test_build_shell_layer_config_accepts_legacy_fallback_keys():
|
||||
agent_soul = AgentSoulConfig.model_validate(
|
||||
{
|
||||
"tools": {
|
||||
"cli_tools": [
|
||||
{"label": "node", "install_command": "apt-get install -y nodejs"},
|
||||
{"tool_name": "python", "setup_command": "pip install pytest"},
|
||||
{"install": "apk add git"},
|
||||
{"ignored": True},
|
||||
]
|
||||
},
|
||||
"env": {
|
||||
"variables": [
|
||||
{"key": "PROJECT_NAME", "default": "demo"},
|
||||
{"env_name": "RETRY_COUNT", "value": 3},
|
||||
{"value": "missing-name"},
|
||||
],
|
||||
"secret_refs": [
|
||||
{"variable": "TOKEN", "credential_id": "credential-1"},
|
||||
{"name": "API_KEY", "provider_credential_id": "credential-2"},
|
||||
{"ref": "missing-name"},
|
||||
],
|
||||
},
|
||||
}
|
||||
)
|
||||
|
||||
config = build_shell_layer_config(agent_soul).model_dump(mode="json")
|
||||
|
||||
assert config["cli_tools"] == [
|
||||
{"name": "node", "install_commands": ["apt-get install -y nodejs"]},
|
||||
{"name": "python", "install_commands": ["pip install pytest"]},
|
||||
{"name": None, "install_commands": ["apk add git"]},
|
||||
]
|
||||
assert config["env"] == [
|
||||
{"name": "PROJECT_NAME", "value": "demo"},
|
||||
{"name": "RETRY_COUNT", "value": "3"},
|
||||
]
|
||||
assert config["secret_refs"] == [
|
||||
{"name": "TOKEN", "ref": "credential-1"},
|
||||
{"name": "API_KEY", "ref": "credential-2"},
|
||||
]
|
||||
assert config["sandbox"] is None
|
||||
|
||||
|
||||
def test_builds_workflow_run_request_with_dify_plugin_tools_layer():
|
||||
@@ -381,6 +463,7 @@ def test_empty_declared_outputs_injects_prd_defaults_text_files_json():
|
||||
assert properties["files"]["type"] == "array"
|
||||
# `files` defaults to array<file> → items is a file ref object.
|
||||
assert properties["files"]["items"]["properties"]["file_id"]["type"] == "string"
|
||||
assert properties["files"]["items"]["required"] == ["file_id"]
|
||||
assert properties["json"]["type"] == "object"
|
||||
# Defaults are all required=False so no `required:` key on the schema.
|
||||
assert "required" not in output_layer["json_schema"]
|
||||
|
||||
@@ -167,6 +167,27 @@ def test_publish_validation_rejects_missing_agent_soul_model():
|
||||
)
|
||||
|
||||
|
||||
def test_publish_validation_rejects_duplicate_cli_tool_names():
|
||||
node_job = WorkflowNodeJobConfig.model_validate({})
|
||||
snapshot = _snapshot()
|
||||
snapshot.config_snapshot = AgentSoulConfig(
|
||||
model=AgentSoulModelConfig(
|
||||
plugin_id="langgenius/openai",
|
||||
model_provider="openai",
|
||||
model="gpt-test",
|
||||
),
|
||||
tools={"cli_tools": [{"name": "pytest"}, {"tool_name": "pytest"}]},
|
||||
)
|
||||
session = Mock()
|
||||
session.scalar.side_effect = [_binding(node_job), _agent(), snapshot]
|
||||
|
||||
with pytest.raises(WorkflowAgentNodeValidationError, match="duplicate CLI Tool name pytest"):
|
||||
WorkflowAgentNodeValidator.validate_published_workflow(
|
||||
session=session,
|
||||
workflow=_workflow(_graph([{"source": "start", "target": "agent-node"}])),
|
||||
)
|
||||
|
||||
|
||||
def test_publish_validation_rejects_missing_previous_node():
|
||||
node_job = WorkflowNodeJobConfig.model_validate(
|
||||
{"previous_node_output_refs": [{"node_id": "missing-node", "output": "text"}]}
|
||||
|
||||
@@ -0,0 +1,284 @@
|
||||
"""Unit tests for the Agent App sandbox workspace inspector service.
|
||||
|
||||
These cover session-id resolution from the conversation snapshot and proxying to
|
||||
the backend client, with fakes for the session store and client (no DB / no HTTP).
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from collections.abc import Generator
|
||||
|
||||
import pytest
|
||||
from agenton.compositor import CompositorSessionSnapshot
|
||||
from agenton.compositor.schemas import LayerSessionSnapshot
|
||||
from agenton.layers.base import LifecycleState
|
||||
from sqlalchemy import delete
|
||||
|
||||
from clients.agent_backend.workspace_files_client import (
|
||||
WorkspaceDownloadResult,
|
||||
WorkspaceFileEntry,
|
||||
WorkspaceListResult,
|
||||
WorkspacePreviewResult,
|
||||
)
|
||||
from core.db.session_factory import session_factory
|
||||
from models.agent import AgentRuntimeSession, AgentRuntimeSessionOwnerType, AgentRuntimeSessionStatus
|
||||
from services.agent_app_workspace_service import (
|
||||
AgentAppWorkspaceService,
|
||||
AgentWorkspaceInspectorError,
|
||||
WorkflowAgentWorkspaceService,
|
||||
_default_client_factory,
|
||||
)
|
||||
|
||||
|
||||
def _snapshot(*, shell: bool = True, session_id: str | None = "abc1234") -> CompositorSessionSnapshot:
|
||||
layers = [LayerSessionSnapshot(name="history", lifecycle_state=LifecycleState.SUSPENDED, runtime_state={})]
|
||||
if shell and session_id is not None:
|
||||
layers.append(
|
||||
LayerSessionSnapshot(
|
||||
name="shell",
|
||||
lifecycle_state=LifecycleState.SUSPENDED,
|
||||
runtime_state={"session_id": session_id, "workspace_cwd": f"~/workspace/{session_id}"},
|
||||
)
|
||||
)
|
||||
elif shell:
|
||||
layers.append(LayerSessionSnapshot(name="shell", lifecycle_state=LifecycleState.SUSPENDED, runtime_state={}))
|
||||
return CompositorSessionSnapshot(layers=layers)
|
||||
|
||||
|
||||
class FakeStore:
|
||||
def __init__(self, snapshot: CompositorSessionSnapshot | None) -> None:
|
||||
self._snapshot = snapshot
|
||||
self.scope: tuple[str, str, str] | None = None
|
||||
|
||||
def load_active_snapshot_for_conversation(
|
||||
self, *, tenant_id: str, app_id: str, conversation_id: str
|
||||
) -> CompositorSessionSnapshot | None:
|
||||
self.scope = (tenant_id, app_id, conversation_id)
|
||||
return self._snapshot
|
||||
|
||||
|
||||
class FakeClient:
|
||||
def __init__(self) -> None:
|
||||
self.calls: list[tuple[str, str, str]] = []
|
||||
|
||||
def list_files(self, session_id: str, path: str) -> WorkspaceListResult:
|
||||
self.calls.append(("list", session_id, path))
|
||||
return WorkspaceListResult(
|
||||
path=path, entries=[WorkspaceFileEntry(name="a.txt", type="file", size=1, mtime=1)], truncated=False
|
||||
)
|
||||
|
||||
def preview(self, session_id: str, path: str) -> WorkspacePreviewResult:
|
||||
self.calls.append(("preview", session_id, path))
|
||||
return WorkspacePreviewResult(path=path, size=5, truncated=False, binary=False, text="hello")
|
||||
|
||||
def download(self, session_id: str, path: str) -> WorkspaceDownloadResult:
|
||||
self.calls.append(("download", session_id, path))
|
||||
return WorkspaceDownloadResult(path=path, size=3, truncated=False, content=b"abc")
|
||||
|
||||
|
||||
def _service(
|
||||
snapshot: CompositorSessionSnapshot | None,
|
||||
) -> tuple[AgentAppWorkspaceService, FakeClient, FakeStore]:
|
||||
store = FakeStore(snapshot)
|
||||
client = FakeClient()
|
||||
service = AgentAppWorkspaceService(session_store=store, client_factory=lambda: client) # type: ignore[arg-type]
|
||||
return service, client, store
|
||||
|
||||
|
||||
def test_list_resolves_session_id_and_proxies() -> None:
|
||||
service, client, store = _service(_snapshot(session_id="abc1234"))
|
||||
|
||||
result = service.list_files(tenant_id="t1", app_id="app1", conversation_id="conv1", path="sub")
|
||||
|
||||
assert result.entries[0].name == "a.txt"
|
||||
assert client.calls == [("list", "abc1234", "sub")]
|
||||
assert store.scope == ("t1", "app1", "conv1")
|
||||
|
||||
|
||||
def test_preview_and_download_use_resolved_session() -> None:
|
||||
service, client, _ = _service(_snapshot(session_id="abc1234"))
|
||||
|
||||
preview = service.preview(tenant_id="t", app_id="a", conversation_id="c", path="n.txt")
|
||||
download = service.download(tenant_id="t", app_id="a", conversation_id="c", path="b.bin")
|
||||
|
||||
assert preview.text == "hello"
|
||||
assert download.content == b"abc"
|
||||
assert client.calls == [("preview", "abc1234", "n.txt"), ("download", "abc1234", "b.bin")]
|
||||
|
||||
|
||||
def test_no_active_session_raises_404() -> None:
|
||||
service, client, _ = _service(None)
|
||||
|
||||
with pytest.raises(AgentWorkspaceInspectorError) as exc_info:
|
||||
service.list_files(tenant_id="t", app_id="a", conversation_id="c", path=".")
|
||||
|
||||
assert exc_info.value.code == "no_active_session"
|
||||
assert exc_info.value.status_code == 404
|
||||
assert client.calls == []
|
||||
|
||||
|
||||
def test_snapshot_without_shell_layer_raises_no_sandbox() -> None:
|
||||
service, _, _ = _service(_snapshot(shell=False))
|
||||
|
||||
with pytest.raises(AgentWorkspaceInspectorError) as exc_info:
|
||||
service.list_files(tenant_id="t", app_id="a", conversation_id="c", path=".")
|
||||
|
||||
assert exc_info.value.code == "no_sandbox"
|
||||
assert exc_info.value.status_code == 404
|
||||
|
||||
|
||||
def test_shell_layer_without_session_id_raises_no_sandbox() -> None:
|
||||
service, _, _ = _service(_snapshot(session_id=None))
|
||||
|
||||
with pytest.raises(AgentWorkspaceInspectorError) as exc_info:
|
||||
service.preview(tenant_id="t", app_id="a", conversation_id="c", path="n.txt")
|
||||
|
||||
assert exc_info.value.code == "no_sandbox"
|
||||
|
||||
|
||||
def test_default_client_factory_requires_agent_backend_base_url(monkeypatch: pytest.MonkeyPatch) -> None:
|
||||
monkeypatch.setattr("services.agent_app_workspace_service.dify_config.AGENT_BACKEND_BASE_URL", "")
|
||||
|
||||
with pytest.raises(AgentWorkspaceInspectorError) as exc_info:
|
||||
_default_client_factory()
|
||||
|
||||
assert exc_info.value.code == "inspector_unavailable"
|
||||
assert exc_info.value.status_code == 503
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
def _runtime_session_table() -> Generator[None, None, None]:
|
||||
engine = session_factory.get_session_maker().kw["bind"]
|
||||
AgentRuntimeSession.__table__.create(bind=engine, checkfirst=True)
|
||||
yield
|
||||
with session_factory.create_session() as session:
|
||||
session.execute(delete(AgentRuntimeSession))
|
||||
session.commit()
|
||||
AgentRuntimeSession.__table__.drop(bind=engine, checkfirst=True)
|
||||
|
||||
|
||||
def _insert_workflow_session(
|
||||
*,
|
||||
workflow_run_id: str = "run-1",
|
||||
node_id: str = "node-1",
|
||||
node_execution_id: str = "node-exec-1",
|
||||
binding_id: str = "binding-1",
|
||||
session_id: str = "abc1234",
|
||||
) -> None:
|
||||
with session_factory.create_session() as session:
|
||||
session.add(
|
||||
AgentRuntimeSession(
|
||||
tenant_id="tenant-1",
|
||||
app_id="app-1",
|
||||
owner_type=AgentRuntimeSessionOwnerType.WORKFLOW_RUN,
|
||||
workflow_id="workflow-1",
|
||||
workflow_run_id=workflow_run_id,
|
||||
node_id=node_id,
|
||||
node_execution_id=node_execution_id,
|
||||
binding_id=binding_id,
|
||||
agent_id="agent-1",
|
||||
agent_config_snapshot_id="snapshot-1",
|
||||
backend_run_id="backend-run-1",
|
||||
session_snapshot=_snapshot(session_id=session_id).model_dump_json(),
|
||||
composition_layer_specs="[]",
|
||||
status=AgentRuntimeSessionStatus.ACTIVE,
|
||||
)
|
||||
)
|
||||
session.commit()
|
||||
|
||||
|
||||
@pytest.mark.usefixtures("_runtime_session_table")
|
||||
def test_workflow_workspace_service_resolves_run_node_session_and_proxies() -> None:
|
||||
_insert_workflow_session(session_id="def5678")
|
||||
client = FakeClient()
|
||||
service = WorkflowAgentWorkspaceService(client_factory=lambda: client) # type: ignore[arg-type]
|
||||
|
||||
result = service.list_files(
|
||||
tenant_id="tenant-1",
|
||||
app_id="app-1",
|
||||
workflow_run_id="run-1",
|
||||
node_id="node-1",
|
||||
node_execution_id="node-exec-1",
|
||||
path=".",
|
||||
)
|
||||
|
||||
assert result.entries[0].name == "a.txt"
|
||||
assert client.calls == [("list", "def5678", ".")]
|
||||
|
||||
|
||||
@pytest.mark.usefixtures("_runtime_session_table")
|
||||
def test_workflow_workspace_service_filters_by_node_execution_id() -> None:
|
||||
_insert_workflow_session(node_execution_id="node-exec-1", session_id="abc1234")
|
||||
_insert_workflow_session(node_execution_id="node-exec-2", binding_id="binding-2", session_id="def5678")
|
||||
client = FakeClient()
|
||||
service = WorkflowAgentWorkspaceService(client_factory=lambda: client) # type: ignore[arg-type]
|
||||
|
||||
_ = service.preview(
|
||||
tenant_id="tenant-1",
|
||||
app_id="app-1",
|
||||
workflow_run_id="run-1",
|
||||
node_id="node-1",
|
||||
node_execution_id="node-exec-2",
|
||||
path="out.txt",
|
||||
)
|
||||
|
||||
assert client.calls == [("preview", "def5678", "out.txt")]
|
||||
|
||||
|
||||
@pytest.mark.usefixtures("_runtime_session_table")
|
||||
def test_workflow_workspace_service_download_uses_latest_active_session_when_execution_id_is_omitted() -> None:
|
||||
_insert_workflow_session(node_execution_id="node-exec-1", session_id="abc1234")
|
||||
client = FakeClient()
|
||||
service = WorkflowAgentWorkspaceService(client_factory=lambda: client) # type: ignore[arg-type]
|
||||
|
||||
result = service.download(
|
||||
tenant_id="tenant-1",
|
||||
app_id="app-1",
|
||||
workflow_run_id="run-1",
|
||||
node_id="node-1",
|
||||
node_execution_id=None,
|
||||
path="out.bin",
|
||||
)
|
||||
|
||||
assert result.content == b"abc"
|
||||
assert client.calls == [("download", "abc1234", "out.bin")]
|
||||
|
||||
|
||||
@pytest.mark.usefixtures("_runtime_session_table")
|
||||
def test_workflow_workspace_service_raises_when_no_active_session_exists() -> None:
|
||||
client = FakeClient()
|
||||
service = WorkflowAgentWorkspaceService(client_factory=lambda: client) # type: ignore[arg-type]
|
||||
|
||||
with pytest.raises(AgentWorkspaceInspectorError) as exc_info:
|
||||
service.list_files(
|
||||
tenant_id="tenant-1",
|
||||
app_id="app-1",
|
||||
workflow_run_id="run-1",
|
||||
node_id="node-1",
|
||||
node_execution_id=None,
|
||||
path=".",
|
||||
)
|
||||
|
||||
assert exc_info.value.code == "no_active_session"
|
||||
assert exc_info.value.status_code == 404
|
||||
assert client.calls == []
|
||||
|
||||
|
||||
@pytest.mark.usefixtures("_runtime_session_table")
|
||||
def test_workflow_workspace_service_raises_when_snapshot_has_no_shell_session() -> None:
|
||||
_insert_workflow_session(session_id="")
|
||||
client = FakeClient()
|
||||
service = WorkflowAgentWorkspaceService(client_factory=lambda: client) # type: ignore[arg-type]
|
||||
|
||||
with pytest.raises(AgentWorkspaceInspectorError) as exc_info:
|
||||
service.preview(
|
||||
tenant_id="tenant-1",
|
||||
app_id="app-1",
|
||||
workflow_run_id="run-1",
|
||||
node_id="node-1",
|
||||
node_execution_id=None,
|
||||
path="out.txt",
|
||||
)
|
||||
|
||||
assert exc_info.value.code == "no_sandbox"
|
||||
assert client.calls == []
|
||||
Reference in New Issue
Block a user