mirror of
https://github.com/langgenius/dify.git
synced 2026-09-24 23:22:26 +08:00
feat(agent): Agent Files / agent Cloud storage — api backend (ENG-589) (#37172)
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
789698cddd
commit
a80bba2c35
@@ -13,7 +13,11 @@ from dataclasses import dataclass
|
||||
from typing import Any, Protocol, cast
|
||||
|
||||
from agenton.compositor import CompositorSessionSnapshot
|
||||
from dify_agent.layers.execution_context import DifyExecutionContextLayerConfig
|
||||
from dify_agent.layers.execution_context import (
|
||||
DifyExecutionContextInvokeFrom,
|
||||
DifyExecutionContextLayerConfig,
|
||||
DifyExecutionContextUserFrom,
|
||||
)
|
||||
from dify_agent.protocol import CreateRunRequest
|
||||
|
||||
from clients.agent_backend import (
|
||||
@@ -126,7 +130,10 @@ class AgentAppRuntimeRequestBuilder:
|
||||
conversation_id=context.conversation_id,
|
||||
agent_id=context.agent_id,
|
||||
agent_config_version_id=context.agent_config_snapshot_id,
|
||||
invoke_from="agent_app",
|
||||
# Agent Files §1.3: real Dify access context + agent run mode.
|
||||
user_from=cast(DifyExecutionContextUserFrom, context.dify_context.user_from.value),
|
||||
invoke_from=cast(DifyExecutionContextInvokeFrom, context.dify_context.invoke_from.value),
|
||||
agent_mode="agent_app",
|
||||
),
|
||||
agent_soul_prompt=agent_soul.prompt.system_prompt or None,
|
||||
user_prompt=context.user_query,
|
||||
|
||||
@@ -231,6 +231,20 @@ class RequestRequestUploadFile(BaseModel):
|
||||
mimetype: str
|
||||
|
||||
|
||||
class RequestRequestDownloadFile(BaseModel):
|
||||
"""Request a signed download URL for a workflow file ref (Agent Files §3.1.1).
|
||||
|
||||
``user_from`` / ``invoke_from`` are the flattened Dify file-access context (the
|
||||
dify-agent server reads them from the execution context). ``file`` is a standard
|
||||
file mapping: ``transfer_method`` plus ``reference`` (local_file / tool_file /
|
||||
datasource_file) or ``url`` (remote_url).
|
||||
"""
|
||||
|
||||
user_from: str
|
||||
invoke_from: str
|
||||
file: Mapping[str, Any]
|
||||
|
||||
|
||||
class RequestFetchAppInfo(BaseModel):
|
||||
"""
|
||||
Request to fetch app info
|
||||
|
||||
@@ -479,8 +479,9 @@ class DifyNodeFactory(NodeFactory):
|
||||
if issubclass(node_class, DifyAgentNode):
|
||||
from clients.agent_backend import AgentBackendRunEventAdapter, AgentBackendRunRequestBuilder
|
||||
from clients.agent_backend.factory import create_agent_backend_run_client
|
||||
from core.workflow.nodes.agent_v2.file_tenant_validator import UploadFileTenantValidator
|
||||
from core.workflow.nodes.agent_v2.file_tenant_validator import AgentOutputFileTenantValidator
|
||||
from core.workflow.nodes.agent_v2.output_failure_orchestrator import OutputFailureOrchestrator
|
||||
from core.workflow.nodes.agent_v2.output_file_rebacker import reback_tool_file_output
|
||||
from core.workflow.nodes.agent_v2.output_type_checker import PerOutputTypeChecker
|
||||
from core.workflow.nodes.agent_v2.session_store import WorkflowAgentRuntimeSessionStore
|
||||
|
||||
@@ -496,11 +497,12 @@ class DifyNodeFactory(NodeFactory):
|
||||
fake_scenario=dify_config.AGENT_BACKEND_FAKE_SCENARIO,
|
||||
),
|
||||
"event_adapter": AgentBackendRunEventAdapter(),
|
||||
"output_adapter": WorkflowAgentOutputAdapter(),
|
||||
# Agent Files §4.6: reback file outputs from the ToolFile row so
|
||||
# downstream metadata is authoritative, not sandbox-provided.
|
||||
"output_adapter": WorkflowAgentOutputAdapter(tool_file_rebacker=reback_tool_file_output),
|
||||
# Stage 4 §5/§7: per-output validation + failure orchestration. The
|
||||
# tenant validator queries upload_files so it stays cheap when
|
||||
# outputs contain no file refs.
|
||||
"type_checker": PerOutputTypeChecker(file_validator=UploadFileTenantValidator()),
|
||||
# tenant validator resolves ToolFile (canonical) + UploadFile refs.
|
||||
"type_checker": PerOutputTypeChecker(file_validator=AgentOutputFileTenantValidator()),
|
||||
"failure_orchestrator": OutputFailureOrchestrator(),
|
||||
"session_store": WorkflowAgentRuntimeSessionStore(),
|
||||
}
|
||||
|
||||
@@ -312,6 +312,7 @@ class DifyAgentNode(Node[DifyAgentNodeData]):
|
||||
inputs=inputs,
|
||||
process_data=process_data,
|
||||
metadata=metadata,
|
||||
tenant_id=dify_ctx.tenant_id,
|
||||
)
|
||||
)
|
||||
return
|
||||
@@ -342,6 +343,7 @@ class DifyAgentNode(Node[DifyAgentNodeData]):
|
||||
inputs=inputs,
|
||||
process_data=process_data,
|
||||
metadata=metadata,
|
||||
tenant_id=dify_ctx.tenant_id,
|
||||
)
|
||||
)
|
||||
return
|
||||
|
||||
@@ -1,13 +1,12 @@
|
||||
"""Tenant-scope validator for file refs produced by Agent backend outputs.
|
||||
|
||||
Stage 4 §5.3: every file output the Agent backend produces must resolve to an
|
||||
``upload_files`` row that belongs to the current tenant; cross-tenant file
|
||||
references must never be plumbed downstream. ``PerOutputTypeChecker`` accepts a
|
||||
``FileTenantValidator`` Protocol so unit tests can stub the check without
|
||||
hitting Postgres.
|
||||
|
||||
This module supplies the production implementation that queries the
|
||||
``upload_files`` table via SQLAlchemy.
|
||||
Stage 4 §5.3 / Agent Files §4.6: every file output the Agent backend produces
|
||||
must resolve to a file record owned by the current tenant; cross-tenant file
|
||||
references must never be plumbed downstream. Agent runtime output files are
|
||||
canonically ``ToolFile`` (referenced by a minimal ``{id}``), so this validator
|
||||
checks ``tool_files`` first and falls back to ``upload_files`` for compatibility
|
||||
with older/manual refs. ``PerOutputTypeChecker`` accepts a ``FileTenantValidator``
|
||||
Protocol so unit tests can stub the check without hitting Postgres.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
@@ -19,10 +18,11 @@ from sqlalchemy.exc import DataError, SQLAlchemyError
|
||||
|
||||
from core.db.session_factory import session_factory
|
||||
from models.model import UploadFile
|
||||
from models.tools import ToolFile
|
||||
|
||||
|
||||
class UploadFileTenantValidator:
|
||||
"""Production ``FileTenantValidator`` backed by the ``upload_files`` table.
|
||||
class AgentOutputFileTenantValidator:
|
||||
"""Production ``FileTenantValidator`` backed by ``tool_files`` + ``upload_files``.
|
||||
|
||||
Returns ``False`` (rejects the file) on any pathological input: empty
|
||||
file_id/tenant_id, non-UUID file_id format, DB errors. The Agent backend
|
||||
@@ -40,7 +40,15 @@ class UploadFileTenantValidator:
|
||||
return False
|
||||
try:
|
||||
with session_factory.create_session() as session:
|
||||
owner_tenant_id = session.scalar(select(UploadFile.tenant_id).where(UploadFile.id == file_id))
|
||||
# Agent output files are canonically ToolFile; check it first.
|
||||
tool_owner = session.scalar(select(ToolFile.tenant_id).where(ToolFile.id == file_id))
|
||||
if tool_owner is not None:
|
||||
return tool_owner == tenant_id
|
||||
upload_owner = session.scalar(select(UploadFile.tenant_id).where(UploadFile.id == file_id))
|
||||
except (DataError, SQLAlchemyError):
|
||||
return False
|
||||
return owner_tenant_id == tenant_id
|
||||
return upload_owner == tenant_id
|
||||
|
||||
|
||||
# Back-compat alias for callers/tests that imported the upload-only name.
|
||||
UploadFileTenantValidator = AgentOutputFileTenantValidator
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
from __future__ import annotations
|
||||
|
||||
from collections.abc import Mapping
|
||||
from collections.abc import Callable, Mapping
|
||||
from typing import Any
|
||||
|
||||
from clients.agent_backend import (
|
||||
@@ -21,6 +21,13 @@ from graphon.variables.segments import ArrayFileSegment, FileSegment
|
||||
class WorkflowAgentOutputAdapter:
|
||||
"""Convert terminal Agent backend events into workflow node run results."""
|
||||
|
||||
def __init__(self, *, tool_file_rebacker: Callable[..., File | None] | None = None) -> None:
|
||||
# Agent Files §4.6: resolve a bare ToolFile id into a graphon File whose
|
||||
# metadata comes from the ToolFile row (not the untrusted sandbox payload).
|
||||
# Injected so unit tests can stub it without DB access; None keeps the
|
||||
# legacy payload-only behaviour for non-file or rich-payload outputs.
|
||||
self._tool_file_rebacker = tool_file_rebacker
|
||||
|
||||
def build_success_result(
|
||||
self,
|
||||
*,
|
||||
@@ -28,6 +35,7 @@ class WorkflowAgentOutputAdapter:
|
||||
inputs: dict[str, Any],
|
||||
process_data: dict[str, Any],
|
||||
metadata: dict[str, Any],
|
||||
tenant_id: str | None = None,
|
||||
) -> NodeRunResult:
|
||||
metadata = self._with_terminal_metadata(metadata, event, "succeeded")
|
||||
usage = self._usage_from_metadata(metadata)
|
||||
@@ -35,7 +43,7 @@ class WorkflowAgentOutputAdapter:
|
||||
status=WorkflowNodeExecutionStatus.SUCCEEDED,
|
||||
inputs=inputs,
|
||||
process_data=process_data,
|
||||
outputs=self._normalize_outputs(event.output),
|
||||
outputs=self._normalize_outputs(event.output, tenant_id=tenant_id),
|
||||
metadata=self._build_node_metadata(metadata=metadata, usage=usage),
|
||||
llm_usage=usage or LLMUsage.empty_usage(),
|
||||
)
|
||||
@@ -101,49 +109,93 @@ class WorkflowAgentOutputAdapter:
|
||||
error_type="agent_backend_stream_error",
|
||||
)
|
||||
|
||||
@classmethod
|
||||
def _normalize_outputs(cls, output: Any) -> dict[str, Any]:
|
||||
def _normalize_outputs(self, output: Any, *, tenant_id: str | None) -> dict[str, Any]:
|
||||
if isinstance(output, dict):
|
||||
if cls._is_file_payload(output):
|
||||
return {"file": cls._file_segment_from_payload(output)}
|
||||
return {key: cls._normalize_output_value(value) for key, value in output.items()}
|
||||
if self._is_file_payload(output):
|
||||
file = self._file_from_payload(output, tenant_id=tenant_id)
|
||||
if file is not None:
|
||||
return {"file": FileSegment(value=file)}
|
||||
return {key: self._normalize_output_value(value, tenant_id=tenant_id) for key, value in output.items()}
|
||||
if isinstance(output, str):
|
||||
return {"text": output}
|
||||
return {"result": output}
|
||||
|
||||
@classmethod
|
||||
def _normalize_output_value(cls, value: Any) -> Any:
|
||||
def _normalize_output_value(self, value: Any, *, tenant_id: str | None) -> Any:
|
||||
if isinstance(value, File | FileSegment | ArrayFileSegment):
|
||||
return value
|
||||
if isinstance(value, Mapping):
|
||||
if cls._is_file_payload(value):
|
||||
return cls._file_segment_from_payload(value)
|
||||
return {key: cls._normalize_output_value(item) for key, item in value.items()}
|
||||
if self._is_file_payload(value):
|
||||
file = self._file_from_payload(value, tenant_id=tenant_id)
|
||||
if file is not None:
|
||||
return FileSegment(value=file)
|
||||
# A bare ref that did not resolve to a tenant file: treat as a plain object.
|
||||
return {key: self._normalize_output_value(item, tenant_id=tenant_id) for key, item in value.items()}
|
||||
if isinstance(value, list):
|
||||
if value and all(isinstance(item, Mapping) and cls._is_file_payload(item) for item in value):
|
||||
return ArrayFileSegment(value=[cls._file_from_payload(item) for item in value])
|
||||
return [cls._normalize_output_value(item) for item in value]
|
||||
if value and all(isinstance(item, Mapping) and self._is_file_payload(item) for item in value):
|
||||
files = [self._file_from_payload(item, tenant_id=tenant_id) for item in value]
|
||||
if all(file is not None for file in files):
|
||||
return ArrayFileSegment(value=[file for file in files if file is not None])
|
||||
return [self._normalize_output_value(item, tenant_id=tenant_id) for item in value]
|
||||
return value
|
||||
|
||||
@staticmethod
|
||||
def _is_file_payload(value: Mapping[str, Any]) -> bool:
|
||||
return any(value.get(key) for key in ("file_id", "upload_file_id", "tool_file_id", "url", "remote_url")) and (
|
||||
"filename" in value or "mime_type" in value or "url" in value or "remote_url" in value
|
||||
# Keys a file-output ref may legitimately carry. A dict is treated as a file
|
||||
# ref only if it has an id/url AND every key is one of these — so a bare
|
||||
# ``{"id": "..."}`` (Agent Files §4.6 canonical) is recognized while ordinary
|
||||
# business objects that merely contain an ``id`` field are not.
|
||||
_FILE_FIELD_KEYS: frozenset[str] = frozenset(
|
||||
{
|
||||
"id",
|
||||
"file_id",
|
||||
"upload_file_id",
|
||||
"tool_file_id",
|
||||
"url",
|
||||
"remote_url",
|
||||
"filename",
|
||||
"name",
|
||||
"mime_type",
|
||||
"mimetype",
|
||||
"extension",
|
||||
"size",
|
||||
"type",
|
||||
"file_type",
|
||||
}
|
||||
)
|
||||
|
||||
@classmethod
|
||||
def _is_file_payload(cls, value: Mapping[str, Any]) -> bool:
|
||||
has_ref = any(
|
||||
isinstance(value.get(key), str) and value.get(key)
|
||||
for key in ("id", "file_id", "upload_file_id", "tool_file_id", "url", "remote_url")
|
||||
)
|
||||
return has_ref and all(key in cls._FILE_FIELD_KEYS for key in value)
|
||||
|
||||
@classmethod
|
||||
def _file_segment_from_payload(cls, value: Mapping[str, Any]) -> FileSegment:
|
||||
return FileSegment(value=cls._file_from_payload(value))
|
||||
@staticmethod
|
||||
def _is_rich_payload(value: Mapping[str, Any]) -> bool:
|
||||
"""The payload carries its own metadata, so it can build a File without DB reback."""
|
||||
return any(value.get(key) for key in ("filename", "name", "mime_type", "mimetype", "url", "remote_url"))
|
||||
|
||||
@classmethod
|
||||
def _file_from_payload(cls, value: Mapping[str, Any]) -> File:
|
||||
remote_url = cls._string_value(value.get("remote_url") or value.get("url"))
|
||||
upload_file_id = cls._string_value(value.get("upload_file_id") or value.get("file_id"))
|
||||
tool_file_id = cls._string_value(value.get("tool_file_id"))
|
||||
filename = cls._string_value(value.get("filename") or value.get("name"))
|
||||
mime_type = cls._string_value(value.get("mime_type") or value.get("mimetype"))
|
||||
extension = cls._extension_from_payload(value, filename)
|
||||
file_type = cls._file_type_from_payload(value, mime_type)
|
||||
def _file_from_payload(self, value: Mapping[str, Any], *, tenant_id: str | None) -> File | None:
|
||||
# Canonical Agent output file is a ToolFile referenced by ``id`` (or the
|
||||
# ``tool_file_id`` alias). Reback its metadata authoritatively from the
|
||||
# ToolFile row instead of trusting the sandbox payload.
|
||||
tool_file_id = self._string_value(value.get("tool_file_id") or value.get("id"))
|
||||
remote_url = self._string_value(value.get("remote_url") or value.get("url"))
|
||||
upload_file_id = self._string_value(value.get("upload_file_id") or value.get("file_id"))
|
||||
|
||||
if tool_file_id and self._tool_file_rebacker is not None and tenant_id:
|
||||
rebacked = self._tool_file_rebacker(tenant_id=tenant_id, tool_file_id=tool_file_id)
|
||||
if rebacked is not None:
|
||||
return rebacked
|
||||
|
||||
# No authoritative reback: only build a File from the payload when it
|
||||
# actually carries file metadata; a bare unresolved id is not a file.
|
||||
if not self._is_rich_payload(value):
|
||||
return None
|
||||
|
||||
filename = self._string_value(value.get("filename") or value.get("name"))
|
||||
mime_type = self._string_value(value.get("mime_type") or value.get("mimetype"))
|
||||
extension = self._extension_from_payload(value, filename)
|
||||
file_type = self._file_type_from_payload(value, mime_type)
|
||||
size = value.get("size")
|
||||
if not isinstance(size, int):
|
||||
size = -1
|
||||
|
||||
@@ -0,0 +1,63 @@
|
||||
"""Reback an Agent backend file output (a bare ``ToolFile`` id) into a graphon File.
|
||||
|
||||
Agent Files §4.6: an agent run returns output files referenced only by id
|
||||
(``{"id": "<tool_file_id>"}``). The authoritative ``filename`` / ``mime_type`` /
|
||||
``extension`` / ``size`` come from the ``ToolFile`` row, never from the
|
||||
(untrusted) sandbox payload. This module resolves a tenant-owned ToolFile id
|
||||
into a full graphon ``File`` so downstream workflow consumers get correct,
|
||||
trustworthy metadata.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from mimetypes import guess_extension
|
||||
from uuid import UUID
|
||||
|
||||
from sqlalchemy import select
|
||||
from sqlalchemy.exc import DataError, SQLAlchemyError
|
||||
|
||||
from core.db.session_factory import session_factory
|
||||
from core.workflow.file_reference import build_file_reference
|
||||
from graphon.file import File, FileTransferMethod, get_file_type_by_mime_type
|
||||
from models.tools import ToolFile
|
||||
|
||||
|
||||
def reback_tool_file_output(*, tenant_id: str, tool_file_id: str) -> File | None:
|
||||
"""Build a graphon File from a ToolFile id owned by ``tenant_id``.
|
||||
|
||||
Returns ``None`` when the id is empty/malformed or does not resolve to a
|
||||
ToolFile owned by the tenant (the caller then treats the value as a plain
|
||||
object rather than fabricating a file with empty metadata).
|
||||
"""
|
||||
if not tool_file_id or not tenant_id:
|
||||
return None
|
||||
try:
|
||||
UUID(tool_file_id)
|
||||
except (ValueError, TypeError):
|
||||
return None
|
||||
try:
|
||||
with session_factory.create_session() as session:
|
||||
tool_file = session.scalar(
|
||||
select(ToolFile).where(ToolFile.id == tool_file_id, ToolFile.tenant_id == tenant_id)
|
||||
)
|
||||
except (DataError, SQLAlchemyError):
|
||||
return None
|
||||
if tool_file is None:
|
||||
return None
|
||||
|
||||
mime_type = tool_file.mimetype or ""
|
||||
extension = guess_extension(mime_type) or ".bin"
|
||||
return File(
|
||||
type=get_file_type_by_mime_type(mime_type),
|
||||
transfer_method=FileTransferMethod.TOOL_FILE,
|
||||
remote_url=None,
|
||||
reference=build_file_reference(record_id=str(tool_file.id)),
|
||||
related_id=tool_file.id,
|
||||
filename=tool_file.name,
|
||||
extension=extension,
|
||||
mime_type=mime_type or None,
|
||||
size=tool_file.size,
|
||||
)
|
||||
|
||||
|
||||
__all__ = ["reback_tool_file_output"]
|
||||
@@ -78,9 +78,9 @@ class FileTenantValidator(Protocol):
|
||||
def is_owned_by_tenant(self, *, file_id: str, tenant_id: str) -> bool: ...
|
||||
|
||||
|
||||
# Recognized aliases the Agent backend (or pydantic-ai) may produce for the
|
||||
# canonical file id field. The canonical spec form is ``file_id`` (§5.2).
|
||||
_FILE_ID_KEYS: tuple[str, ...] = ("file_id", "upload_file_id", "tool_file_id")
|
||||
# Recognized id fields in a file-output ref. Agent Files §4.6: the canonical
|
||||
# minimal form is ``{"id": "<tool_file_id>"}``; the rest are accepted aliases.
|
||||
_FILE_ID_KEYS: tuple[str, ...] = ("id", "file_id", "upload_file_id", "tool_file_id")
|
||||
|
||||
|
||||
class PerOutputTypeChecker:
|
||||
|
||||
@@ -5,7 +5,11 @@ from dataclasses import dataclass
|
||||
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.execution_context import (
|
||||
DifyExecutionContextInvokeFrom,
|
||||
DifyExecutionContextLayerConfig,
|
||||
DifyExecutionContextUserFrom,
|
||||
)
|
||||
from dify_agent.layers.shell import (
|
||||
DifyShellCliToolConfig,
|
||||
DifyShellEnvVarConfig,
|
||||
@@ -178,7 +182,12 @@ class WorkflowAgentRuntimeRequestBuilder:
|
||||
conversation_id=get_system_text(context.variable_pool, SystemVariableKey.CONVERSATION_ID),
|
||||
agent_id=context.agent.id,
|
||||
agent_config_version_id=context.snapshot.id,
|
||||
invoke_from=self._agent_backend_invoke_from(context.dify_context.invoke_from),
|
||||
# Agent Files §1.3: forward the real Dify access context
|
||||
# (user_from + invoke_from) so downstream file/drive inner APIs
|
||||
# can rebuild it; the agent run mode moves to agent_mode.
|
||||
user_from=cast(DifyExecutionContextUserFrom, context.dify_context.user_from.value),
|
||||
invoke_from=cast(DifyExecutionContextInvokeFrom, context.dify_context.invoke_from.value),
|
||||
agent_mode=self._agent_mode(context.dify_context.invoke_from),
|
||||
),
|
||||
agent_soul_prompt=agent_soul.prompt.system_prompt or None,
|
||||
workflow_node_job_prompt=workflow_job_prompt,
|
||||
@@ -202,7 +211,7 @@ class WorkflowAgentRuntimeRequestBuilder:
|
||||
)
|
||||
|
||||
@staticmethod
|
||||
def _agent_backend_invoke_from(invoke_from: InvokeFrom) -> Literal["workflow_run", "single_step"]:
|
||||
def _agent_mode(invoke_from: InvokeFrom) -> Literal["workflow_run", "single_step"]:
|
||||
if invoke_from in {InvokeFrom.DEBUGGER, InvokeFrom.VALIDATION}:
|
||||
return "single_step"
|
||||
return "workflow_run"
|
||||
|
||||
Reference in New Issue
Block a user