chore(agent-v2): sync daily changes (#38162)

Co-authored-by: yunlu.wen <yunlu.wen@dify.ai>
Co-authored-by: autofix-ci[bot] <114827586+autofix-ci[bot]@users.noreply.github.com>
Co-authored-by: Yunlu Wen <wylswz@163.com>
Co-authored-by: Copilot Autofix powered by AI <175728472+Copilot@users.noreply.github.com>
Co-authored-by: Joel <iamjoel007@gmail.com>
Co-authored-by: Yanli 盐粒 <yanli@dify.ai>
Co-authored-by: 盐粒 Yanli <beautyyuyanli@gmail.com>
Co-authored-by: zyssyz123 <916125788@qq.com>
Co-authored-by: 盐粒 Yanli <mail@yanli.one>
This commit is contained in:
yyh
2026-07-01 05:07:23 +00:00
committed by GitHub
co-authored by yunlu.wen autofix-ci[bot] Yunlu Wen Copilot Autofix powered by AI Joel Yanli 盐粒 盐粒 Yanli zyssyz123 盐粒 Yanli
parent f816ae2e95
commit 0923ebaf88
277 changed files with 26866 additions and 7348 deletions
+170 -32
View File
@@ -1,21 +1,28 @@
"""Agent App generator: orchestrate one conversation turn for an Agent App.
"""Agent App generator: orchestrate Agent App chat and finalize executions.
Mirrors the agent_chat generator (conversation + message + queue + streamed
response over the EasyUI chat pipeline), but the backing config comes from the
bound Agent Soul and the answer is produced by ``AgentAppRunner`` calling the
dify-agent backend rather than an in-process LLM/ReAct loop.
The primary mode mirrors the agent_chat generator (conversation + message +
queue + streamed response over the EasyUI chat pipeline), but the backing
config comes from the bound Agent Soul and the answer is produced by
``AgentAppRunner`` calling the dify-agent backend rather than an in-process
LLM/ReAct loop.
It also exposes a stateless build-finalize mode that reuses existing runtime
context from the bound debug conversation, triggers the Agent backend side
effect synchronously, and skips Dify-side chat/message persistence.
"""
from __future__ import annotations
import contextvars
import json
import logging
import threading
import uuid
from collections.abc import Generator, Mapping
from typing import Any
from collections.abc import Generator, Mapping, Sequence
from typing import Any, Literal
from flask import Flask, current_app
from pydantic import JsonValue
from sqlalchemy import and_, or_, select
from clients.agent_backend import AgentBackendRunEventAdapter
@@ -61,6 +68,13 @@ class AgentAppGeneratorError(ValueError):
"""Raised when an Agent App turn cannot be set up."""
def _append_prompt_file_mappings(query: str, prompt_file_mappings: Sequence[JsonValue]) -> str:
"""Append raw request file references to the backend user prompt."""
if not prompt_file_mappings:
return query
return f"{query}\n{json.dumps(list(prompt_file_mappings), ensure_ascii=False)}"
class AgentAppGenerator(MessageBasedAppGenerator):
def generate(
self,
@@ -74,14 +88,12 @@ class AgentAppGenerator(MessageBasedAppGenerator):
if not streaming:
raise AgentAppGeneratorError("Agent App only supports streaming mode")
query = args.get("query")
if not isinstance(query, str) or not query.strip():
raise AgentAppGeneratorError("query is required")
query = query.replace("\x00", "")
query = self._require_query(args)
inputs = args["inputs"]
prompt_file_mappings = args.get("files") or []
# Resolve the bound roster Agent + its current Agent Soul snapshot.
agent, agent_config_id, agent_soul = self._resolve_agent(
agent, agent_config_id, agent_config_version_kind, agent_soul = self._resolve_agent(
app_model,
invoke_from=invoke_from,
draft_type=args.get("draft_type"),
@@ -122,6 +134,7 @@ class AgentAppGenerator(MessageBasedAppGenerator):
),
query=query,
files=[],
prompt_file_mappings=prompt_file_mappings,
parent_message_id=(
args.get("parent_message_id")
if invoke_from not in {InvokeFrom.SERVICE_API, InvokeFrom.OPENAPI}
@@ -137,6 +150,7 @@ class AgentAppGenerator(MessageBasedAppGenerator):
trace_manager=trace_manager,
agent_id=agent.id,
agent_config_snapshot_id=agent_config_id,
agent_config_version_kind=agent_config_version_kind,
agent_runtime_session_snapshot_id=runtime_session_snapshot_id,
)
@@ -176,6 +190,86 @@ class AgentAppGenerator(MessageBasedAppGenerator):
)
return AgentAppGenerateResponseConverter.convert(response=response, invoke_from=invoke_from)
def generate_stateless(
self,
*,
app_model: App,
user: Account | EndUser,
args: Mapping[str, Any],
invoke_from: InvokeFrom,
) -> Mapping[str, Any]:
"""Run one Agent App turn without persisting Dify conversation messages."""
query = self._require_query(args)
conversation_id = args.get("conversation_id")
if not isinstance(conversation_id, str) or not conversation_id:
raise AgentAppGeneratorError("conversation_id is required")
agent, agent_config_id, agent_config_version_kind, agent_soul = self._resolve_agent(
app_model,
invoke_from=invoke_from,
draft_type=args.get("draft_type"),
user=user,
)
runtime_session_snapshot_id = self._runtime_session_snapshot_id(
invoke_from=invoke_from,
snapshot_id=agent_config_id,
)
return self._run_stateless(
app_model=app_model,
user=user,
invoke_from=invoke_from,
query=query,
conversation_id=conversation_id,
agent=agent,
agent_config_id=agent_config_id,
agent_config_version_kind=agent_config_version_kind,
agent_soul=agent_soul,
runtime_session_snapshot_id=runtime_session_snapshot_id,
)
def _run_stateless(
self,
*,
app_model: App,
user: Account | EndUser,
invoke_from: InvokeFrom,
query: str,
conversation_id: str,
agent: Agent,
agent_config_id: str,
agent_config_version_kind: Literal["snapshot", "draft", "build_draft"],
agent_soul: AgentSoulConfig,
runtime_session_snapshot_id: str | None,
) -> Mapping[str, Any]:
"""Run the Agent backend without creating or updating Dify chat records.
Build-chat finalization is an action against the Agent backend (for
example, ``dify-agent config push``). It may reuse the active build-chat
runtime snapshot for shell/config context, but the API side must not add
a synthetic user/assistant turn to the debug conversation.
"""
dify_context = DifyRunContext(
tenant_id=app_model.tenant_id,
app_id=app_model.id,
user_id=user.id,
user_from=UserFrom.ACCOUNT if isinstance(user, Account) else UserFrom.END_USER,
invoke_from=invoke_from,
)
self._build_runner(dify_context).run_stateless(
dify_context=dify_context,
agent_id=agent.id,
agent_config_snapshot_id=agent_config_id,
agent_config_version_kind=agent_config_version_kind,
agent_soul=agent_soul,
conversation_id=conversation_id,
query=query,
idempotency_key=str(uuid.uuid4()),
session_scope_snapshot_id=runtime_session_snapshot_id,
)
return {"result": "success"}
def resume_after_form_submission(
self,
*,
@@ -192,15 +286,15 @@ class AgentAppGenerator(MessageBasedAppGenerator):
persisted to the conversation. Live streaming to a reconnected client is
out of scope here — the message is persisted and can be re-fetched.
"""
agent, agent_config_id, agent_soul = self._resolve_agent(
app_model,
invoke_from=invoke_from,
draft_type="draft",
user=user,
)
conversation = ConversationService.get_conversation(
app_model=app_model, conversation_id=conversation_id, user=user
)
agent, agent_config_id, agent_config_version_kind, agent_soul = self._resolve_agent(
app_model,
invoke_from=invoke_from,
draft_type=self._resume_draft_type(app_model=app_model, conversation=conversation, user=user),
user=user,
)
app_config = AgentAppConfigManager.get_app_config(
app_model=app_model,
@@ -245,6 +339,7 @@ class AgentAppGenerator(MessageBasedAppGenerator):
trace_manager=trace_manager,
agent_id=agent.id,
agent_config_snapshot_id=agent_config_id,
agent_config_version_kind=agent_config_version_kind,
)
conversation, message = self._init_generate_records(application_generate_entity, conversation)
@@ -285,6 +380,30 @@ class AgentAppGenerator(MessageBasedAppGenerator):
stream=False,
)
@staticmethod
def _resume_draft_type(*, app_model: App, conversation: Any, user: Account | EndUser) -> str | None:
if conversation.invoke_from != InvokeFrom.DEBUGGER:
return None
active_session = AgentAppRuntimeSessionStore().load_active_session_for_conversation(
tenant_id=app_model.tenant_id,
app_id=app_model.id,
conversation_id=conversation.id,
)
snapshot_id = active_session.scope.agent_config_snapshot_id if active_session is not None else None
if snapshot_id and isinstance(user, Account):
draft = db.session.scalar(
select(AgentConfigDraft).where(
AgentConfigDraft.tenant_id == app_model.tenant_id,
AgentConfigDraft.id == snapshot_id,
)
)
if draft is not None:
if draft.draft_type == AgentConfigDraftType.DEBUG_BUILD and draft.account_id == user.id:
return AgentConfigDraftType.DEBUG_BUILD.value
if draft.draft_type == AgentConfigDraftType.DRAFT and draft.account_id is None:
return AgentConfigDraftType.DRAFT.value
return AgentConfigDraftType.DRAFT.value
def _generate_worker(
self,
*,
@@ -329,6 +448,10 @@ class AgentAppGenerator(MessageBasedAppGenerator):
)
if handled:
return
query = _append_prompt_file_mappings(
query=query,
prompt_file_mappings=application_generate_entity.prompt_file_mappings,
)
dify_context = DifyRunContext(
tenant_id=app_config.tenant_id,
@@ -337,27 +460,18 @@ class AgentAppGenerator(MessageBasedAppGenerator):
user_from=user_from,
invoke_from=application_generate_entity.invoke_from,
)
credentials_provider, _ = build_dify_model_access(dify_context)
_, _, agent_soul = self._resolve_agent_by_id(
tenant_id=app_config.tenant_id,
agent_id=application_generate_entity.agent_id,
snapshot_id=application_generate_entity.agent_config_snapshot_id,
)
runner = AgentAppRunner(
request_builder=AgentAppRuntimeRequestBuilder(credentials_provider=credentials_provider),
agent_backend_client=create_agent_backend_run_client(
base_url=dify_config.AGENT_BACKEND_BASE_URL,
use_fake=dify_config.AGENT_BACKEND_USE_FAKE,
fake_scenario=dify_config.AGENT_BACKEND_FAKE_SCENARIO,
),
event_adapter=AgentBackendRunEventAdapter(),
session_store=AgentAppRuntimeSessionStore(),
)
runner = self._build_runner(dify_context)
runner.run(
dify_context=dify_context,
agent_id=application_generate_entity.agent_id,
agent_config_snapshot_id=application_generate_entity.agent_config_snapshot_id,
agent_config_version_kind=application_generate_entity.agent_config_version_kind,
agent_soul=agent_soul,
conversation_id=conversation.id,
query=query,
@@ -374,6 +488,27 @@ class AgentAppGenerator(MessageBasedAppGenerator):
finally:
db.session.close()
@staticmethod
def _require_query(args: Mapping[str, Any]) -> str:
query = args.get("query")
if not isinstance(query, str) or not query.strip():
raise AgentAppGeneratorError("query is required")
return query.replace("\x00", "")
@staticmethod
def _build_runner(dify_context: DifyRunContext) -> AgentAppRunner:
credentials_provider, _ = build_dify_model_access(dify_context)
return AgentAppRunner(
request_builder=AgentAppRuntimeRequestBuilder(credentials_provider=credentials_provider),
agent_backend_client=create_agent_backend_run_client(
base_url=dify_config.AGENT_BACKEND_BASE_URL,
use_fake=dify_config.AGENT_BACKEND_USE_FAKE,
fake_scenario=dify_config.AGENT_BACKEND_FAKE_SCENARIO,
),
event_adapter=AgentBackendRunEventAdapter(),
session_store=AgentAppRuntimeSessionStore(),
)
def _run_input_guards(
self,
*,
@@ -446,7 +581,7 @@ class AgentAppGenerator(MessageBasedAppGenerator):
invoke_from: InvokeFrom,
draft_type: Any,
user: Account | EndUser,
) -> tuple[Agent, str, AgentSoulConfig]:
) -> tuple[Agent, str, Literal["snapshot", "draft", "build_draft"], AgentSoulConfig]:
agent = db.session.scalar(
select(Agent)
.where(
@@ -474,13 +609,16 @@ class AgentAppGenerator(MessageBasedAppGenerator):
account_id=user.id if isinstance(user, Account) else None,
)
agent_soul = AgentSoulConfig.model_validate(draft.config_snapshot_dict)
return agent, draft.id, agent_soul
config_version_kind: Literal["snapshot", "draft", "build_draft"] = (
"build_draft" if draft.draft_type == AgentConfigDraftType.DEBUG_BUILD else "draft"
)
return agent, draft.id, config_version_kind, agent_soul
_, snapshot, agent_soul = self._resolve_agent_by_id(
tenant_id=app_model.tenant_id,
agent_id=agent.id,
snapshot_id=agent.active_config_snapshot_id,
)
return agent, snapshot.id, agent_soul
return agent, snapshot.id, "snapshot", agent_soul
@staticmethod
def _runtime_session_snapshot_id(*, invoke_from: InvokeFrom, snapshot_id: str) -> str | None:
+436 -34
View File
@@ -1,18 +1,23 @@
"""Agent App runner: drive one conversation turn through the dify-agent backend.
"""Agent App runner: drive Agent backend turns for both chat and finalize flows.
Unlike the legacy ``AgentChatAppRunner`` (which runs an in-process ReAct loop),
this runner delegates to the Agent backend: build the run request from the
Agent Soul + conversation, create the run, consume its event stream, and
republish the assistant answer as chat queue events so the existing
EasyUI chat task pipeline persists the message and streams SSE. The conversation
``session_snapshot`` is saved on success for multi-turn continuity (S3).
this runner delegates to the Agent backend and supports two execution modes.
- Normal chat turns build the run request from the Agent Soul + conversation,
consume backend stream events, republish the assistant answer through the
existing EasyUI chat task pipeline, and save the conversation
``session_snapshot`` on success for multi-turn continuity (S3).
- Stateless build-finalize turns reuse any prior conversation snapshot only to
construct the backend request, wait synchronously for backend completion, and
intentionally do not persist Dify-side chat records or runtime-session state.
"""
from __future__ import annotations
import json
import logging
from typing import Any
from decimal import Decimal
from typing import Any, Literal
from dify_agent.layers.ask_human import AskHumanToolArgs
from dify_agent.protocol import DeferredToolResultsPayload
@@ -28,6 +33,7 @@ from clients.agent_backend import (
AgentBackendStreamInternalEvent,
extract_runtime_layer_specs,
)
from configs import dify_config
from core.app.apps.agent_app.runtime_request_builder import (
AgentAppRuntimeBuildContext,
AgentAppRuntimeRequest,
@@ -41,13 +47,16 @@ from core.app.apps.agent_app.session_store import (
from core.app.apps.base_app_queue_manager import AppQueueManager, PublishFrom
from core.app.apps.exc import GenerateTaskStoppedError
from core.app.entities.app_invoke_entities import DifyRunContext
from core.app.entities.queue_entities import QueueLLMChunkEvent, QueueMessageEndEvent
from core.app.entities.queue_entities import QueueAgentThoughtEvent, QueueLLMChunkEvent, QueueMessageEndEvent
from core.repositories.human_input_repository import HumanInputFormRepository, HumanInputFormRepositoryImpl
from core.workflow.nodes.agent_v2.ask_human_hitl import AskHumanFormBuildError, create_ask_human_form
from core.workflow.nodes.agent_v2.ask_human_resume import build_deferred_tool_results, resolve_ask_human_form
from extensions.ext_database import db
from graphon.model_runtime.entities.llm_entities import LLMResult, LLMResultChunk, LLMResultChunkDelta, LLMUsage
from graphon.model_runtime.entities.message_entities import AssistantPromptMessage, PromptMessage, UserPromptMessage
from models.agent_config_entities import AgentSoulConfig
from models.enums import CreatorUserRole
from models.model import MessageAgentThought
logger = logging.getLogger(__name__)
@@ -134,6 +143,283 @@ def publish_message_end(
)
class _AgentProcessRecorder:
"""Persist Agent v2 thinking/tool process events through the legacy thought model."""
def __init__(
self,
*,
dify_context: DifyRunContext,
message_id: str,
queue_manager: AppQueueManager,
) -> None:
self._dify_context = dify_context
self._message_id = message_id
self._queue_manager = queue_manager
self._next_position = 1
self._thinking_by_index: dict[int, str] = {}
self._tool_by_index: dict[int, str] = {}
self._tool_by_call_id: dict[str, str] = {}
self._open_tool_by_name: dict[str, set[str]] = {}
def handle_stream_event(self, event: AgentBackendStreamInternalEvent) -> None:
data = event.data
if not isinstance(data, dict):
return
event_kind = data.get("event_kind")
if event_kind == "part_delta":
self._handle_part_delta(data)
elif event_kind == "part_start":
self._handle_part(data)
elif event_kind in {"function_tool_call", "output_tool_call"}:
self._handle_tool_call_event(data)
elif event_kind in {"function_tool_result", "output_tool_result"}:
self._handle_tool_result_event(data)
def _handle_part_delta(self, data: dict[str, Any]) -> None:
delta = data.get("delta")
if not isinstance(delta, dict):
return
index = _event_index(data)
delta_kind = delta.get("part_delta_kind")
if delta_kind == "thinking":
content_delta = delta.get("content_delta")
if isinstance(content_delta, str) and content_delta:
self._append_thinking(index, content_delta)
return
if delta_kind == "tool_call":
self._record_tool_call_delta(index, delta)
def _handle_part(self, data: dict[str, Any]) -> None:
part = data.get("part")
if not isinstance(part, dict):
return
index = _event_index(data)
part_kind = part.get("part_kind")
if part_kind == "thinking":
content = part.get("content")
if isinstance(content, str) and content:
self._append_thinking(index, content)
return
if part_kind in {"tool-call", "builtin-tool-call"}:
self._record_tool_call_part(index, part)
return
if part_kind in {"tool-return", "builtin-tool-return"}:
self._record_tool_return_part(part)
def _handle_tool_call_event(self, data: dict[str, Any]) -> None:
part = data.get("part")
if isinstance(part, dict):
self._record_tool_call_part(_event_index(data), part)
def _handle_tool_result_event(self, data: dict[str, Any]) -> None:
part = data.get("part") or data.get("result")
if isinstance(part, dict):
self._record_tool_return_part(part)
return
content = data.get("content")
if content is not None:
self._record_tool_observation(
tool_call_id=_string_or_none(data.get("tool_call_id")),
tool_name=_string_or_none(data.get("tool_name")),
observation=content,
)
def _append_thinking(self, index: int, content_delta: str) -> None:
thought_id = self._thinking_by_index.get(index)
if thought_id is None:
thought_id = self._create_thought(thought=content_delta)
self._thinking_by_index[index] = thought_id
return
self._update_thought(thought_id, thought_delta=content_delta)
def _record_tool_call_delta(self, index: int, delta: dict[str, Any]) -> None:
tool_call_id = _string_or_none(delta.get("tool_call_id"))
tool_name = _string_or_none(delta.get("tool_name_delta"))
args_delta = delta.get("args_delta")
thought_id = self._lookup_tool_thought(index=index, tool_call_id=tool_call_id)
if thought_id is None:
thought_id = self._create_thought(tool=tool_name, tool_input=_json_or_text(args_delta))
self._remember_tool_thought(
index=index, tool_call_id=tool_call_id, tool_name=tool_name, thought_id=thought_id
)
return
self._update_thought(
thought_id,
tool=tool_name,
tool_input_delta=_json_or_text(args_delta),
)
def _record_tool_call_part(self, index: int, part: dict[str, Any]) -> None:
tool_call_id = _string_or_none(part.get("tool_call_id"))
tool_name = _string_or_none(part.get("tool_name"))
thought_id = self._lookup_tool_thought(index=index, tool_call_id=tool_call_id)
if thought_id is None:
thought_id = self._create_thought(tool=tool_name, tool_input=_json_or_text(part.get("args")))
self._remember_tool_thought(
index=index, tool_call_id=tool_call_id, tool_name=tool_name, thought_id=thought_id
)
return
self._update_thought(
thought_id,
tool=tool_name,
tool_input=_json_or_text(part.get("args")),
)
def _record_tool_return_part(self, part: dict[str, Any]) -> None:
tool_call_id = _string_or_none(part.get("tool_call_id"))
tool_name = _string_or_none(part.get("tool_name"))
content = part.get("content")
if content is None:
content = part
self._record_tool_observation(tool_call_id=tool_call_id, tool_name=tool_name, observation=content)
def _record_tool_observation(self, *, tool_call_id: str | None, tool_name: str | None, observation: Any) -> None:
thought_id = self._lookup_observation_thought(tool_call_id=tool_call_id, tool_name=tool_name)
if thought_id is None:
thought_id = self._create_thought(tool=tool_name)
else:
self._mark_tool_observed(thought_id)
self._update_thought(thought_id, observation=_json_or_text(observation))
def _lookup_tool_thought(self, *, index: int, tool_call_id: str | None) -> str | None:
if tool_call_id and tool_call_id in self._tool_by_call_id:
return self._tool_by_call_id[tool_call_id]
return self._tool_by_index.get(index)
def _remember_tool_thought(
self, *, index: int, tool_call_id: str | None, tool_name: str | None, thought_id: str
) -> None:
self._tool_by_index[index] = thought_id
if tool_call_id:
self._tool_by_call_id[tool_call_id] = thought_id
if tool_name:
self._open_tool_by_name.setdefault(tool_name, set()).add(thought_id)
def _lookup_observation_thought(self, *, tool_call_id: str | None, tool_name: str | None) -> str | None:
if tool_call_id:
return self._tool_by_call_id.get(tool_call_id)
if tool_name:
open_thought_ids = self._open_tool_by_name.get(tool_name, set())
if len(open_thought_ids) == 1:
return next(iter(open_thought_ids))
return None
def _mark_tool_observed(self, thought_id: str) -> None:
for open_thought_ids in self._open_tool_by_name.values():
open_thought_ids.discard(thought_id)
def _create_thought(
self, *, thought: str | None = None, tool: str | None = None, tool_input: str | None = None
) -> str:
row = MessageAgentThought(
message_id=self._message_id,
message_chain_id=None,
thought=thought,
tool=tool,
tool_labels_str=_tool_labels(tool),
tool_meta_str="{}",
tool_input=tool_input,
observation=None,
tool_process_data=None,
message=None,
message_token=0,
message_unit_price=Decimal(0),
message_price_unit=Decimal("0.001"),
message_files="",
answer="",
answer_token=0,
answer_unit_price=Decimal(0),
answer_price_unit=Decimal("0.001"),
tokens=0,
total_price=Decimal(0),
position=self._next_position,
currency="USD",
latency=0,
created_by_role=self._created_by_role(),
created_by=self._dify_context.user_id,
)
self._next_position += 1
db.session.add(row)
db.session.commit()
thought_id = str(row.id)
self._queue_manager.publish(
QueueAgentThoughtEvent(agent_thought_id=thought_id), PublishFrom.APPLICATION_MANAGER
)
return thought_id
def _update_thought(
self,
thought_id: str,
*,
thought_delta: str | None = None,
tool: str | None = None,
tool_input: str | None = None,
tool_input_delta: str | None = None,
observation: str | None = None,
) -> None:
row = db.session.get(MessageAgentThought, thought_id)
if row is None:
return
if thought_delta:
row.thought = f"{row.thought or ''}{thought_delta}"
if tool:
row.tool = tool
row.tool_labels_str = _tool_labels(tool)
if tool_input is not None:
row.tool_input = tool_input
if tool_input_delta:
row.tool_input = f"{row.tool_input or ''}{tool_input_delta}"
if observation is not None:
row.observation = observation
db.session.commit()
self._queue_manager.publish(
QueueAgentThoughtEvent(agent_thought_id=thought_id), PublishFrom.APPLICATION_MANAGER
)
def _created_by_role(self) -> CreatorUserRole:
if self._dify_context.invoke_from.runs_as_account():
return CreatorUserRole.ACCOUNT
return CreatorUserRole.END_USER
def _event_index(data: dict[str, Any]) -> int:
index = data.get("index")
return index if isinstance(index, int) else -1
def _string_or_none(value: Any) -> str | None:
return value if isinstance(value, str) and value else None
def _json_or_text(value: Any) -> str | None:
if value is None:
return None
if isinstance(value, str):
return value
try:
return json.dumps(value, ensure_ascii=False)
except Exception:
return str(value)
def _tool_labels(tool: str | None) -> str:
if not tool:
return "{}"
return json.dumps({tool: {"en_US": tool, "zh_Hans": tool}}, ensure_ascii=False)
class AgentAppRunner:
"""Runs one Agent App conversation turn against the Agent backend."""
@@ -156,6 +442,7 @@ class AgentAppRunner:
dify_context: DifyRunContext,
agent_id: str,
agent_config_snapshot_id: str,
agent_config_version_kind: Literal["snapshot", "draft", "build_draft"] = "snapshot",
agent_soul: AgentSoulConfig,
conversation_id: str,
query: str,
@@ -164,42 +451,34 @@ class AgentAppRunner:
queue_manager: AppQueueManager,
session_scope_snapshot_id: str | None | _DefaultSessionScopeSnapshotId = _DEFAULT_SESSION_SCOPE_SNAPSHOT_ID,
) -> None:
if isinstance(session_scope_snapshot_id, _DefaultSessionScopeSnapshotId):
effective_session_scope_snapshot_id: str | None = agent_config_snapshot_id
else:
effective_session_scope_snapshot_id = session_scope_snapshot_id
scope = AgentAppSessionScope(
tenant_id=dify_context.tenant_id,
app_id=dify_context.app_id,
conversation_id=conversation_id,
scope = self._build_session_scope(
dify_context=dify_context,
agent_id=agent_id,
agent_config_snapshot_id=effective_session_scope_snapshot_id,
agent_config_snapshot_id=agent_config_snapshot_id,
conversation_id=conversation_id,
session_scope_snapshot_id=session_scope_snapshot_id,
)
# ENG-638: if a prior turn paused on ask_human and the form is now answered,
# resume by threading the human's reply into this run as deferred_tool_results.
stored = self._session_store.load_active_session(scope)
session_snapshot = stored.session_snapshot if stored is not None else None
deferred_tool_results = self._resolve_pending_ask_human(
stored=stored, dify_context=dify_context, message_id=message_id
)
runtime = self._request_builder.build(
AgentAppRuntimeBuildContext(
dify_context=dify_context,
agent_id=agent_id,
agent_config_snapshot_id=agent_config_snapshot_id,
agent_soul=agent_soul,
conversation_id=conversation_id,
user_query=query,
idempotency_key=message_id,
session_snapshot=session_snapshot,
deferred_tool_results=deferred_tool_results,
)
runtime = self._build_runtime(
dify_context=dify_context,
agent_id=agent_id,
agent_config_snapshot_id=agent_config_snapshot_id,
agent_config_version_kind=agent_config_version_kind,
agent_soul=agent_soul,
conversation_id=conversation_id,
query=query,
idempotency_key=message_id,
stored=stored,
message_id=message_id,
)
create_response = self._agent_backend_client.create_run(runtime.request)
terminal, streamed_answer = self._consume_stream(
create_response.run_id,
dify_context=dify_context,
message_id=message_id,
queue_manager=queue_manager,
model_name=model_name,
query=query,
@@ -241,6 +520,111 @@ class AgentAppRunner:
runtime_layer_specs=extract_runtime_layer_specs(runtime.request.composition),
)
def run_stateless(
self,
*,
dify_context: DifyRunContext,
agent_id: str,
agent_config_snapshot_id: str,
agent_config_version_kind: Literal["snapshot", "draft", "build_draft"] = "snapshot",
agent_soul: AgentSoulConfig,
conversation_id: str,
query: str,
idempotency_key: str,
session_scope_snapshot_id: str | None | _DefaultSessionScopeSnapshotId = _DEFAULT_SESSION_SCOPE_SNAPSHOT_ID,
) -> None:
"""Run the Agent backend without creating Dify chat message records.
This path is used by build-chat finalization: the API must trigger the
backend side effects in the existing conversation session, but it must
not persist a synthetic user/assistant turn, update API-side runtime
session rows, or set up HITL state that depends on one.
"""
scope = self._build_session_scope(
dify_context=dify_context,
agent_id=agent_id,
agent_config_snapshot_id=agent_config_snapshot_id,
conversation_id=conversation_id,
session_scope_snapshot_id=session_scope_snapshot_id,
)
runtime = self._build_runtime(
dify_context=dify_context,
agent_id=agent_id,
agent_config_snapshot_id=agent_config_snapshot_id,
agent_config_version_kind=agent_config_version_kind,
agent_soul=agent_soul,
conversation_id=conversation_id,
query=query,
idempotency_key=idempotency_key,
stored=self._session_store.load_active_session(scope),
message_id=None,
)
create_response = self._agent_backend_client.create_run(runtime.request)
status = self._agent_backend_client.wait_run(
create_response.run_id,
timeout_seconds=dify_config.APP_MAX_EXECUTION_TIME,
)
if status.status != "succeeded":
error = getattr(status, "error", None) or f"Agent backend run ended with status {status.status}."
raise AgentBackendError(str(error))
def _build_session_scope(
self,
*,
dify_context: DifyRunContext,
agent_id: str,
agent_config_snapshot_id: str,
conversation_id: str,
session_scope_snapshot_id: str | None | _DefaultSessionScopeSnapshotId,
) -> AgentAppSessionScope:
if isinstance(session_scope_snapshot_id, _DefaultSessionScopeSnapshotId):
effective_session_scope_snapshot_id: str | None = agent_config_snapshot_id
else:
effective_session_scope_snapshot_id = session_scope_snapshot_id
return AgentAppSessionScope(
tenant_id=dify_context.tenant_id,
app_id=dify_context.app_id,
conversation_id=conversation_id,
agent_id=agent_id,
agent_config_snapshot_id=effective_session_scope_snapshot_id,
)
def _build_runtime(
self,
*,
dify_context: DifyRunContext,
agent_id: str,
agent_config_snapshot_id: str,
agent_config_version_kind: Literal["snapshot", "draft", "build_draft"],
agent_soul: AgentSoulConfig,
conversation_id: str,
query: str,
idempotency_key: str,
stored: StoredAgentAppSession | None,
message_id: str | None,
) -> AgentAppRuntimeRequest:
session_snapshot = stored.session_snapshot if stored is not None else None
deferred_tool_results = (
self._resolve_pending_ask_human(stored=stored, dify_context=dify_context, message_id=message_id)
if message_id is not None
else None
)
return self._request_builder.build(
AgentAppRuntimeBuildContext(
dify_context=dify_context,
agent_id=agent_id,
agent_config_snapshot_id=agent_config_snapshot_id,
agent_config_version_kind=agent_config_version_kind,
agent_soul=agent_soul,
conversation_id=conversation_id,
user_query=query,
idempotency_key=idempotency_key,
session_snapshot=session_snapshot,
deferred_tool_results=deferred_tool_results,
)
)
def _pause_for_ask_human(
self,
*,
@@ -334,12 +718,19 @@ class AgentAppRunner:
self,
run_id: str,
*,
dify_context: DifyRunContext,
message_id: str,
queue_manager: AppQueueManager,
model_name: str,
query: str | None,
):
terminal = None
streamed_answer_parts: list[str] = []
process_recorder = _AgentProcessRecorder(
dify_context=dify_context,
message_id=message_id,
queue_manager=queue_manager,
)
for public_event in self._agent_backend_client.stream_events(run_id):
if queue_manager.is_stopped():
self._cancel_run(run_id)
@@ -353,6 +744,17 @@ class AgentAppRunner:
AgentBackendInternalEventType.STREAM_EVENT,
):
if isinstance(internal_event, AgentBackendStreamInternalEvent):
try:
process_recorder.handle_stream_event(internal_event)
except Exception:
db.session.rollback()
logger.warning(
"Failed to persist Agent App process event: run_id=%s message_id=%s event_kind=%s",
run_id,
message_id,
internal_event.event_kind,
exc_info=True,
)
text_delta = self._extract_stream_text_delta(internal_event)
if text_delta:
streamed_answer_parts.append(text_delta)
@@ -12,7 +12,7 @@ from __future__ import annotations
from collections.abc import Mapping
from dataclasses import dataclass
from typing import Any, Protocol, cast
from typing import Any, Literal, Protocol, cast
from agenton.compositor import CompositorSessionSnapshot
from dify_agent.layers.execution_context import (
@@ -29,20 +29,22 @@ from clients.agent_backend import (
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.app.entities.app_invoke_entities import DifyRunContext, InvokeFrom
from core.workflow.nodes.agent_v2.dify_tools_builder import (
WorkflowAgentDifyToolLayersBuilder,
WorkflowAgentDifyToolsBuilder,
WorkflowAgentDifyToolsBuildError,
WorkflowAgentToolLayers,
)
from core.workflow.nodes.agent_v2.runtime_request_builder import (
append_runtime_warnings,
build_ask_human_layer_config,
build_drive_aware_soul_mention_resolver,
build_drive_layer_config,
build_config_aware_soul_mention_resolver,
build_config_layer_config,
build_knowledge_layer_config,
build_shell_layer_config,
)
from models.agent_config_entities import AgentSoulConfig
from models.agent_config_entities import AgentSoulConfig, AgentSoulToolsConfig
from models.provider_ids import ModelProviderID
from services.agent.prompt_mentions import build_soul_mention_resolver, expand_prompt_mentions
@@ -68,6 +70,7 @@ class AgentAppRuntimeBuildContext:
conversation_id: str
user_query: str
idempotency_key: str
agent_config_version_kind: Literal["snapshot", "draft", "build_draft"] = "snapshot"
session_snapshot: CompositorSessionSnapshot | None = None
# ENG-638: set when resuming a chat turn after a submitted ask_human form.
deferred_tool_results: DeferredToolResultsPayload | None = None
@@ -88,11 +91,11 @@ class AgentAppRuntimeRequestBuilder:
*,
credentials_provider: CredentialsProvider,
request_builder: AgentBackendRunRequestBuilder | None = None,
plugin_tools_builder: WorkflowAgentPluginToolsBuilder | None = None,
dify_tools_builder: WorkflowAgentDifyToolLayersBuilder | None = None,
) -> None:
self._credentials_provider = credentials_provider
self._request_builder = request_builder or AgentBackendRunRequestBuilder()
self._plugin_tools_builder = plugin_tools_builder or WorkflowAgentPluginToolsBuilder()
self._dify_tools_builder = dify_tools_builder or WorkflowAgentDifyToolsBuilder()
def build(self, context: AgentAppRuntimeBuildContext) -> AgentAppRuntimeRequest:
agent_soul = context.agent_soul
@@ -105,38 +108,33 @@ class AgentAppRuntimeRequestBuilder:
metadata = self._build_metadata(context)
credentials = self._credentials_provider.fetch(agent_soul.model.model_provider, agent_soul.model.model)
try:
tools_layer = self._plugin_tools_builder.build(
tool_layers = self._build_tool_layers(
tenant_id=context.dify_context.tenant_id,
app_id=context.dify_context.app_id,
user_id=context.dify_context.user_id,
tools=agent_soul.tools,
invoke_from=context.dify_context.invoke_from,
)
except WorkflowAgentPluginToolsBuildError as error:
except WorkflowAgentDifyToolsBuildError as error:
raise AgentAppRuntimeRequestBuildError(error.error_code, str(error)) from error
if tools_layer is not None or agent_soul.tools.cli_tools:
if tool_layers.plugin_tools is not None or tool_layers.core_tools is not None or agent_soul.tools.cli_tools:
metadata["agent_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 [],
"dify_tool_count": len(tool_layers.exposed_tool_names()),
"dify_tool_names": tool_layers.exposed_tool_names(),
"cli_tool_count": len(agent_soul.tools.cli_tools),
}
drive_config = None
config_layer_config = None
soul_prompt_resolver = build_soul_mention_resolver(agent_soul)
if dify_config.AGENT_DRIVE_MANIFEST_ENABLED:
drive_config, drive_warnings = build_drive_layer_config(
config_layer_config, config_warnings = build_config_layer_config(
agent_soul,
tenant_id=context.dify_context.tenant_id,
agent_id=context.agent_id,
)
append_runtime_warnings(metadata, drive_warnings)
soul_prompt_resolver = build_drive_aware_soul_mention_resolver(
agent_soul,
tenant_id=context.dify_context.tenant_id,
agent_id=context.agent_id,
config_version_id=context.agent_config_snapshot_id,
config_version_kind=context.agent_config_version_kind,
)
append_runtime_warnings(metadata, config_warnings)
soul_prompt_resolver = build_config_aware_soul_mention_resolver(agent_soul)
knowledge_config = build_knowledge_layer_config(agent_soul)
request = self._request_builder.build_for_agent_app(
@@ -158,6 +156,7 @@ class AgentAppRuntimeRequestBuilder:
conversation_id=context.conversation_id,
agent_id=context.agent_id,
agent_config_version_id=context.agent_config_snapshot_id,
agent_config_version_kind=context.agent_config_version_kind,
# 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),
@@ -168,9 +167,10 @@ class AgentAppRuntimeRequestBuilder:
agent_soul_prompt=expand_prompt_mentions(agent_soul.prompt.system_prompt, soul_prompt_resolver).strip()
or None,
user_prompt=context.user_query,
tools=tools_layer,
tools=tool_layers.plugin_tools,
core_tools=tool_layers.core_tools,
knowledge=knowledge_config,
drive_config=drive_config,
config_layer_config=config_layer_config,
ask_human_config=build_ask_human_layer_config(agent_soul),
include_shell=dify_config.AGENT_SHELL_ENABLED,
shell_config=build_shell_layer_config(agent_soul),
@@ -183,6 +183,26 @@ class AgentAppRuntimeRequestBuilder:
redacted = cast(dict[str, Any], redact_for_agent_backend_log(request))
return AgentAppRuntimeRequest(request=request, redacted_request=redacted, metadata=metadata)
def _build_tool_layers(
self,
*,
tenant_id: str,
app_id: str,
user_id: str | None,
tools: AgentSoulToolsConfig,
invoke_from: InvokeFrom,
) -> WorkflowAgentToolLayers:
# Production Agent App runs intentionally keep existing plugin configs
# on the direct `dify.plugin.tools` route. This builder emits plugin
# tools directly and non-plugin Dify tools through `dify.core.tools`.
return self._dify_tools_builder.build_layers(
tenant_id=tenant_id,
app_id=app_id,
user_id=user_id,
tools=tools,
invoke_from=invoke_from,
)
@staticmethod
def _build_metadata(context: AgentAppRuntimeBuildContext) -> dict[str, Any]:
return {
+15 -2
View File
@@ -1,8 +1,8 @@
from collections.abc import Mapping, Sequence
from enum import StrEnum
from typing import TYPE_CHECKING, Any
from typing import TYPE_CHECKING, Any, Literal
from pydantic import BaseModel, ConfigDict, Field, ValidationInfo, field_validator
from pydantic import BaseModel, ConfigDict, Field, JsonValue, ValidationInfo, field_validator
from constants import UUID_NIL
from core.app.app_config.entities import EasyUIBasedAppConfig, WorkflowUIBasedAppConfig
@@ -220,11 +220,24 @@ class AgentAppGenerateEntity(ChatAppGenerateEntity):
accepted-entity union. The answer is produced by the dify-agent backend
rather than an in-process LLM call; ``model_conf`` is synthesized from the
bound Agent Soul model so the chat task pipeline can persist usage.
``agent_config_version_kind`` selects which Agent config surface the
backend should read from: immutable snapshot, shared draft, or per-user
build draft.
``agent_runtime_session_snapshot_id`` carries the runtime session scope
used to resume or suspend within the same editable config surface.
``prompt_file_mappings`` preserves the raw request ``files`` array for the
Agent backend prompt. These references are appended to the backend prompt
text while the stored chat message keeps the user's original query.
"""
agent_id: str
agent_config_snapshot_id: str
agent_config_version_kind: Literal["snapshot", "draft", "build_draft"] = "snapshot"
agent_runtime_session_snapshot_id: str | None = None
prompt_file_mappings: Sequence[JsonValue] = Field(default_factory=list)
class AdvancedChatAppGenerateEntity(ConversationAppGenerateEntity):
+3 -5
View File
@@ -111,7 +111,7 @@ class ToolEngine:
tool_messages=binary_files, agent_message=message, invoke_from=invoke_from, user_id=user_id
)
plain_text = ToolEngine._convert_tool_response_to_str(message_list)
plain_text = ToolEngine.tool_response_to_str(message_list)
meta = invocation_meta_dict["meta"]
@@ -234,10 +234,8 @@ class ToolEngine:
yield meta
@staticmethod
def _convert_tool_response_to_str(tool_response: list[ToolInvokeMessage]) -> str:
"""
Handle tool response
"""
def tool_response_to_str(tool_response: list[ToolInvokeMessage]) -> str:
"""Convert tool invoke messages into the plain-text observation shown to the model/user."""
parts: list[str] = []
json_parts: list[str] = []
@@ -0,0 +1,500 @@
from __future__ import annotations
from collections.abc import Mapping
from dataclasses import dataclass
from typing import Any, Literal, Protocol, cast
from dify_agent.layers.dify_core_tools import DifyCoreToolConfig, DifyCoreToolProviderType, DifyCoreToolsLayerConfig
from dify_agent.layers.dify_plugin import (
DifyPluginCredentialValue,
DifyPluginToolConfig,
DifyPluginToolCredentialType,
DifyPluginToolParameter,
DifyPluginToolParameterForm,
DifyPluginToolsLayerConfig,
)
from sqlalchemy import select
from sqlalchemy.orm import Session
from core.agent.entities import AgentToolEntity
from core.app.entities.app_invoke_entities import InvokeFrom
from core.tools.__base.tool import Tool
from core.tools.entities.tool_entities import ToolProviderType
from core.tools.errors import ToolProviderCredentialValidationError, ToolProviderNotFoundError
from core.tools.tool_manager import ToolManager
from core.tools.workflow_as_tool.provider import WorkflowToolProviderController
from extensions.ext_database import db
from models.agent_config_entities import AgentSoulDifyToolConfig, AgentSoulToolsConfig
from models.provider_ids import ToolProviderID
from models.tools import WorkflowToolProvider
from services.tools.mcp_tools_manage_service import MCPToolManageService
class WorkflowAgentDifyToolsBuildError(ValueError):
"""Raised when Agent Soul tools cannot be prepared for Agent backend."""
def __init__(self, error_code: str, message: str) -> None:
self.error_code = error_code
super().__init__(message)
class AgentToolRuntimeProvider(Protocol):
def get_agent_tool_runtime(
self,
tenant_id: str,
app_id: str,
agent_tool: AgentToolEntity,
user_id: str | None = None,
invoke_from: InvokeFrom = InvokeFrom.DEBUGGER,
variable_pool: Any | None = None,
allow_file_parameters: bool = False,
use_default_for_missing_form_parameters: bool = False,
) -> Tool: ...
class ProviderToolsLister(Protocol):
def __call__(
self,
*,
tenant_id: str,
provider_type: ToolProviderType,
provider_id: str,
) -> list[str]: ...
class MCPProviderIDResolver(Protocol):
def __call__(self, *, tenant_id: str, provider_id: str) -> str: ...
@dataclass(frozen=True, slots=True)
class WorkflowAgentToolLayers:
plugin_tools: DifyPluginToolsLayerConfig | None = None
core_tools: DifyCoreToolsLayerConfig | None = None
def exposed_tool_names(self) -> list[str]:
names: list[str] = []
if self.plugin_tools is not None:
names.extend(tool.name or tool.tool_name for tool in self.plugin_tools.tools)
if self.core_tools is not None:
names.extend(tool.name or tool.tool_name for tool in self.core_tools.tools)
return names
class WorkflowAgentDifyToolLayersBuilder(Protocol):
def build_layers(
self,
*,
tenant_id: str,
app_id: str,
user_id: str | None,
tools: AgentSoulToolsConfig,
invoke_from: InvokeFrom,
) -> WorkflowAgentToolLayers: ...
def _list_provider_tool_names(
*,
tenant_id: str,
provider_type: ToolProviderType,
provider_id: str,
) -> list[str]:
"""Tool names a provider currently declares for provider-level Agent entries."""
match provider_type:
case ToolProviderType.PLUGIN:
plugin_provider = ToolManager.get_plugin_provider(provider_id, tenant_id)
return [tool.entity.identity.name for tool in plugin_provider.get_tools() or []]
case ToolProviderType.BUILT_IN:
builtin_provider = ToolManager.get_builtin_provider(provider_id, tenant_id)
return [tool.entity.identity.name for tool in builtin_provider.get_tools() or []]
case ToolProviderType.API:
api_provider, _ = ToolManager.get_api_provider_controller(tenant_id, provider_id)
return [tool.entity.identity.name for tool in api_provider.get_tools(tenant_id) or []]
case ToolProviderType.WORKFLOW:
db_provider = db.session.scalar(
select(WorkflowToolProvider)
.where(
WorkflowToolProvider.id == provider_id,
WorkflowToolProvider.tenant_id == tenant_id,
)
.limit(1)
)
if db_provider is None:
raise ToolProviderNotFoundError(f"workflow provider {provider_id} not found")
workflow_provider = WorkflowToolProviderController.from_db(db_provider)
return [tool.entity.identity.name for tool in workflow_provider.get_tools(tenant_id) or []]
case ToolProviderType.MCP:
mcp_provider = ToolManager.get_mcp_provider_controller(tenant_id, provider_id)
return [tool.entity.identity.name for tool in mcp_provider.get_tools() or []]
case _:
raise ToolProviderNotFoundError(f"provider type {provider_type.value} not found")
def _resolve_mcp_provider_id(*, tenant_id: str, provider_id: str) -> str:
"""Normalize MCP provider ids to the runtime-facing server identifier."""
service = MCPToolManageService(session=cast(Session, db.session))
try:
return service.get_provider_entity(provider_id, tenant_id, by_server_id=True).provider_id
except ValueError:
try:
return service.get_provider_entity(provider_id, tenant_id, by_server_id=False).provider_id
except ValueError as exc:
raise ToolProviderNotFoundError(f"mcp provider {provider_id} not found") from exc
class WorkflowAgentDifyToolsBuilder:
"""Prepare Agent Soul Dify tools for Agent backend run-layer configs.
Plugin tools keep their existing direct daemon path. Core-routed tools
(`builtin`/`api`/`workflow`/`mcp`) are emitted as `dify.core.tools`.
"""
def __init__(
self,
*,
tool_runtime_provider: AgentToolRuntimeProvider | None = None,
provider_tools_lister: ProviderToolsLister | None = None,
mcp_provider_id_resolver: MCPProviderIDResolver | None = None,
) -> None:
self._tool_runtime_provider = tool_runtime_provider or ToolManager
self._provider_tools_lister = provider_tools_lister or _list_provider_tool_names
self._mcp_provider_id_resolver = mcp_provider_id_resolver or _resolve_mcp_provider_id
def build_layers(
self,
*,
tenant_id: str,
app_id: str,
user_id: str | None,
tools: AgentSoulToolsConfig,
invoke_from: InvokeFrom,
) -> WorkflowAgentToolLayers:
"""Resolve user-selected Dify tools into direct/core Agent backend DTOs.
`invoke_from` is the real runtime caller category (DEBUGGER for a
Composer test run, SERVICE_API / WEB_APP for a published run). It must
be threaded through to `ToolManager` so credential quotas, rate limits,
and audit tags match the actual call site.
"""
enabled_tools = [tool for tool in tools.dify_tools if tool.enabled]
if not enabled_tools:
return WorkflowAgentToolLayers()
prepared_plugin: list[DifyPluginToolConfig] = []
prepared_core: list[DifyCoreToolConfig] = []
seen_names: set[str] = set()
for tool_config in self._expand_provider_entries(tenant_id=tenant_id, enabled_tools=enabled_tools):
normalized_tool_config = self._normalized_tool_config(tenant_id=tenant_id, tool_config=tool_config)
destination = self._tool_layer_destination(normalized_tool_config)
exposed_name = self._exposed_tool_name(normalized_tool_config)
if exposed_name in seen_names:
raise WorkflowAgentDifyToolsBuildError(
"agent_tool_name_duplicated",
f"Duplicate Dify Tool name {exposed_name!r}.",
)
seen_names.add(exposed_name)
agent_tool = self._to_agent_tool_entity(normalized_tool_config)
tool_runtime = self._fetch_tool_runtime(
tenant_id=tenant_id,
app_id=app_id,
user_id=user_id,
agent_tool=agent_tool,
invoke_from=invoke_from,
tool_config=normalized_tool_config,
)
if destination == "plugin":
prepared_plugin.append(
self._to_plugin_backend_tool_config(normalized_tool_config, tool_runtime, exposed_name)
)
else:
prepared_core.append(
self._to_core_backend_tool_config(normalized_tool_config, tool_runtime, exposed_name)
)
return WorkflowAgentToolLayers(
plugin_tools=DifyPluginToolsLayerConfig(tools=prepared_plugin) if prepared_plugin else None,
core_tools=DifyCoreToolsLayerConfig(tools=prepared_core) if prepared_core else None,
)
def _expand_provider_entries(
self,
*,
tenant_id: str,
enabled_tools: list[AgentSoulDifyToolConfig],
) -> list[AgentSoulDifyToolConfig]:
"""Expand provider-level entries (`tool_name` omitted = all tools)."""
explicit_by_provider: dict[tuple[ToolProviderType, str], set[str]] = {}
for tool_config in enabled_tools:
if tool_config.tool_name is not None:
explicit_by_provider.setdefault(self._provider_key(tool_config), set()).add(tool_config.tool_name)
expanded: list[AgentSoulDifyToolConfig] = []
for tool_config in enabled_tools:
if tool_config.tool_name is not None:
expanded.append(tool_config)
continue
provider_type = ToolProviderType.value_of(tool_config.provider_type)
provider_id = self._provider_id(tool_config)
try:
tool_names = self._provider_declared_tool_names(
tenant_id=tenant_id,
provider_type=provider_type,
provider_id=provider_id,
)
except ToolProviderNotFoundError as exc:
raise WorkflowAgentDifyToolsBuildError(
"agent_tool_declaration_not_found",
f"Dify Tool provider {provider_id!r} declaration not found: {exc}",
) from exc
if not tool_names:
raise WorkflowAgentDifyToolsBuildError(
"agent_tool_declaration_not_found",
f"Dify Tool provider {provider_id!r} declares no tools.",
)
already_explicit = explicit_by_provider.get(self._provider_key(tool_config), set())
for tool_name in tool_names:
if tool_name in already_explicit:
continue
expanded.append(tool_config.model_copy(update={"tool_name": tool_name, "runtime_parameters": {}}))
return expanded
def _provider_declared_tool_names(
self,
*,
tenant_id: str,
provider_type: ToolProviderType,
provider_id: str,
) -> list[str]:
return self._provider_tools_lister(
tenant_id=tenant_id,
provider_type=provider_type,
provider_id=provider_id,
)
def _normalized_tool_config(
self,
*,
tenant_id: str,
tool_config: AgentSoulDifyToolConfig,
) -> AgentSoulDifyToolConfig:
if tool_config.provider_type != ToolProviderType.MCP.value:
return tool_config
provider_id = self._mcp_provider_id_resolver(tenant_id=tenant_id, provider_id=self._provider_id(tool_config))
return tool_config.model_copy(update={"provider_id": provider_id, "plugin_id": None, "provider": None})
def _fetch_tool_runtime(
self,
*,
tenant_id: str,
app_id: str,
user_id: str | None,
agent_tool: AgentToolEntity,
invoke_from: InvokeFrom,
tool_config: AgentSoulDifyToolConfig,
) -> Tool:
"""Resolve the API-side `Tool` runtime and map fetch errors to stable codes."""
try:
return self._tool_runtime_provider.get_agent_tool_runtime(
tenant_id=tenant_id,
app_id=app_id,
agent_tool=agent_tool,
user_id=user_id,
invoke_from=invoke_from,
variable_pool=None,
allow_file_parameters=True,
use_default_for_missing_form_parameters=True,
)
except ToolProviderNotFoundError as exc:
raise WorkflowAgentDifyToolsBuildError(
"agent_tool_declaration_not_found",
f"Dify Tool {tool_config.tool_name!r} declaration not found: {exc}",
) from exc
except ToolProviderCredentialValidationError as exc:
raise WorkflowAgentDifyToolsBuildError(
"agent_tool_credential_invalid",
f"Dify Tool {tool_config.tool_name!r} credential validation failed: {exc}",
) from exc
except ValueError as exc:
raise WorkflowAgentDifyToolsBuildError(
"agent_tool_config_invalid",
f"Dify Tool {tool_config.tool_name!r} runtime construction failed: {exc}",
) from exc
@staticmethod
def _to_agent_tool_entity(tool_config: AgentSoulDifyToolConfig) -> AgentToolEntity:
assert tool_config.tool_name is not None
return AgentToolEntity(
provider_type=ToolProviderType.value_of(tool_config.provider_type),
provider_id=WorkflowAgentDifyToolsBuilder._provider_id(tool_config),
tool_name=tool_config.tool_name,
tool_parameters=dict(tool_config.runtime_parameters),
credential_id=tool_config.credential_ref.id if tool_config.credential_ref else None,
)
@staticmethod
def _provider_id(tool_config: AgentSoulDifyToolConfig) -> str:
if tool_config.provider_id:
return tool_config.provider_id
assert tool_config.plugin_id is not None
assert tool_config.provider is not None
return f"{tool_config.plugin_id}/{tool_config.provider}"
@staticmethod
def _provider_key(tool_config: AgentSoulDifyToolConfig) -> tuple[ToolProviderType, str]:
return (
ToolProviderType.value_of(tool_config.provider_type),
WorkflowAgentDifyToolsBuilder._provider_id(tool_config),
)
@staticmethod
def _tool_layer_destination(tool_config: AgentSoulDifyToolConfig) -> Literal["plugin", "core"]:
provider_type = ToolProviderType.value_of(tool_config.provider_type)
if provider_type is ToolProviderType.PLUGIN:
return "plugin"
if provider_type in {
ToolProviderType.BUILT_IN,
ToolProviderType.API,
ToolProviderType.WORKFLOW,
ToolProviderType.MCP,
}:
return "core"
if provider_type is ToolProviderType.DATASET_RETRIEVAL:
raise WorkflowAgentDifyToolsBuildError(
"agent_tool_provider_not_supported",
"dataset-retrieval remains on the knowledge path and is not supported in Agent tool layers.",
)
raise WorkflowAgentDifyToolsBuildError(
"agent_tool_provider_not_supported",
f"Dify Tool provider type {provider_type.value!r} is not supported in Agent tool layers.",
)
@staticmethod
def _exposed_tool_name(tool_config: AgentSoulDifyToolConfig) -> str:
assert tool_config.tool_name is not None
return tool_config.tool_name
def _to_plugin_backend_tool_config(
self,
tool_config: AgentSoulDifyToolConfig,
tool_runtime: Tool,
exposed_name: str,
) -> DifyPluginToolConfig:
runtime = tool_runtime.runtime
if runtime is None:
raise WorkflowAgentDifyToolsBuildError(
"agent_tool_config_invalid",
f"Dify Tool {tool_config.tool_name!r} has no runtime.",
)
provider_id = self._provider_id(tool_config)
plugin_id, provider = self._plugin_provider(tool_config, provider_id)
parameters = self._prepared_parameters(tool_runtime)
runtime_parameters = self._runtime_parameters(tool_runtime, parameters)
description = self._description(tool_config, tool_runtime)
return DifyPluginToolConfig(
plugin_id=plugin_id,
provider=provider,
tool_name=exposed_name,
credential_type=self._credential_type(tool_config, runtime.credentials),
name=exposed_name,
description=description,
credentials=self._normalize_credentials(runtime.credentials, tool_name=exposed_name),
runtime_parameters=runtime_parameters,
parameters=parameters,
parameters_json_schema=tool_runtime.get_llm_parameters_json_schema(),
)
def _to_core_backend_tool_config(
self,
tool_config: AgentSoulDifyToolConfig,
tool_runtime: Tool,
exposed_name: str,
) -> DifyCoreToolConfig:
parameters = self._prepared_parameters(tool_runtime)
return DifyCoreToolConfig(
provider_type=cast(DifyCoreToolProviderType, tool_config.provider_type),
provider_id=self._provider_id(tool_config),
tool_name=tool_config.tool_name or exposed_name,
credential_id=tool_config.credential_ref.id if tool_config.credential_ref else None,
name=exposed_name,
description=self._description(tool_config, tool_runtime),
runtime_parameters=self._runtime_parameters(tool_runtime, parameters),
parameters=parameters,
parameters_json_schema=tool_runtime.get_llm_parameters_json_schema(),
)
@staticmethod
def _plugin_provider(tool_config: AgentSoulDifyToolConfig, provider_id: str) -> tuple[str, str]:
if tool_config.plugin_id and tool_config.provider:
return tool_config.plugin_id, tool_config.provider
provider_id_entity = ToolProviderID(provider_id)
return provider_id_entity.plugin_id, provider_id_entity.provider_name
@staticmethod
def _credential_type(
tool_config: AgentSoulDifyToolConfig,
credentials: Mapping[str, Any],
) -> DifyPluginToolCredentialType:
if not credentials and tool_config.credential_type == "unauthorized":
return "unauthorized"
return tool_config.credential_type
@staticmethod
def _prepared_parameters(tool_runtime: Tool) -> list[DifyPluginToolParameter]:
return [
DifyPluginToolParameter.model_validate(parameter.model_dump(mode="json"))
for parameter in tool_runtime.get_merged_runtime_parameters()
]
@staticmethod
def _description(tool_config: AgentSoulDifyToolConfig, tool_runtime: Tool) -> str | None:
description = tool_config.description
if description is None and tool_runtime.entity.description is not None:
description = tool_runtime.entity.description.llm
return description
@staticmethod
def _runtime_parameters(
tool_runtime: Tool,
parameters: list[DifyPluginToolParameter],
) -> dict[str, Any]:
runtime = tool_runtime.runtime
runtime_parameters = dict(runtime.runtime_parameters if runtime is not None else {})
missing = [
parameter.name
for parameter in parameters
if parameter.form is not DifyPluginToolParameterForm.LLM
and parameter.required
and parameter.default is None
and parameter.name not in runtime_parameters
]
if missing:
names = ", ".join(sorted(missing))
raise WorkflowAgentDifyToolsBuildError(
"agent_tool_runtime_parameter_missing",
f"Dify Tool {tool_runtime.entity.identity.name!r} is missing runtime parameters: {names}.",
)
return runtime_parameters
@staticmethod
def _normalize_credentials(
credentials: Mapping[str, Any],
*,
tool_name: str,
) -> dict[str, DifyPluginCredentialValue]:
normalized: dict[str, DifyPluginCredentialValue] = {}
for key, value in credentials.items():
if isinstance(value, str | int | float | bool) or value is None:
normalized[key] = value
continue
raise WorkflowAgentDifyToolsBuildError(
"agent_tool_credential_shape_invalid",
(
f"Dify Plugin Tool {tool_name!r} credential {key!r} has a non-scalar value "
f"({type(value).__name__}); only str/int/float/bool/None are forwarded to the daemon."
),
)
return normalized
@@ -1,334 +0,0 @@
from __future__ import annotations
from collections.abc import Mapping
from typing import Any, Protocol
from dify_agent.layers.dify_plugin import (
DifyPluginCredentialValue,
DifyPluginToolConfig,
DifyPluginToolCredentialType,
DifyPluginToolParameter,
DifyPluginToolParameterForm,
DifyPluginToolsLayerConfig,
)
from core.agent.entities import AgentToolEntity
from core.app.entities.app_invoke_entities import InvokeFrom
from core.tools.__base.tool import Tool
from core.tools.entities.tool_entities import ToolProviderType
from core.tools.errors import (
ToolProviderCredentialValidationError,
ToolProviderNotFoundError,
)
from core.tools.tool_manager import ToolManager
from models.agent_config_entities import AgentSoulDifyToolConfig, AgentSoulToolsConfig
from models.provider_ids import ToolProviderID
class WorkflowAgentPluginToolsBuildError(ValueError):
"""Raised when Agent Soul tools cannot be prepared for Agent backend."""
def __init__(self, error_code: str, message: str) -> None:
self.error_code = error_code
super().__init__(message)
class AgentToolRuntimeProvider(Protocol):
def get_agent_tool_runtime(
self,
tenant_id: str,
app_id: str,
agent_tool: AgentToolEntity,
user_id: str | None = None,
invoke_from: InvokeFrom = InvokeFrom.DEBUGGER,
variable_pool: Any | None = None,
allow_file_parameters: bool = False,
use_default_for_missing_form_parameters: bool = False,
) -> Tool: ...
class ProviderToolsLister(Protocol):
def __call__(self, *, tenant_id: str, provider_id: str) -> list[str]: ...
def _list_provider_tool_names(*, tenant_id: str, provider_id: str) -> list[str]:
"""Tool names a provider currently declares (provider-level config entries)."""
provider = ToolManager.get_builtin_provider(provider_id, tenant_id)
return [tool.entity.identity.name for tool in provider.get_tools() or []]
class WorkflowAgentPluginToolsBuilder:
"""Prepare Agent Soul Dify Plugin Tools for the public Agent backend DTO."""
def __init__(
self,
*,
tool_runtime_provider: AgentToolRuntimeProvider | None = None,
provider_tools_lister: ProviderToolsLister | None = None,
) -> None:
self._tool_runtime_provider = tool_runtime_provider or ToolManager
self._provider_tools_lister = provider_tools_lister or _list_provider_tool_names
def build(
self,
*,
tenant_id: str,
app_id: str,
user_id: str | None,
tools: AgentSoulToolsConfig,
invoke_from: InvokeFrom,
) -> DifyPluginToolsLayerConfig | None:
"""Resolve user-selected Dify Plugin Tools into the Agent backend DTO.
``invoke_from`` is the *real* runtime caller category (DEBUGGER for a
Composer test run, SERVICE_API / WEB_APP for a published run). It must
be threaded through to :class:`ToolManager` so credential quotas, rate
limits, and audit tags match the actual call site.
"""
enabled_tools = [tool for tool in tools.dify_tools if tool.enabled]
if not enabled_tools:
return None
prepared: list[DifyPluginToolConfig] = []
seen_names: set[str] = set()
for tool_config in self._expand_provider_entries(tenant_id=tenant_id, enabled_tools=enabled_tools):
agent_tool = self._to_agent_tool_entity(tool_config)
tool_runtime = self._fetch_tool_runtime(
tenant_id=tenant_id,
app_id=app_id,
user_id=user_id,
agent_tool=agent_tool,
invoke_from=invoke_from,
tool_config=tool_config,
)
exposed_name = self._exposed_tool_name(tool_config)
if exposed_name in seen_names:
raise WorkflowAgentPluginToolsBuildError(
"agent_tool_name_duplicated",
f"Duplicate Dify Plugin Tool name {exposed_name!r}.",
)
seen_names.add(exposed_name)
prepared.append(self._to_backend_tool_config(tool_config, tool_runtime, exposed_name))
return DifyPluginToolsLayerConfig(tools=prepared)
def _expand_provider_entries(
self,
*,
tenant_id: str,
enabled_tools: list[AgentSoulDifyToolConfig],
) -> list[AgentSoulDifyToolConfig]:
"""Expand provider-level entries (``tool_name`` omitted = all tools).
An explicit per-tool entry of the same provider wins over the expansion
(it may carry its own ``runtime_parameters``); expanded clones share the
provider entry's ``credential_ref`` and start with default parameters.
"""
explicit_by_provider: dict[str, set[str]] = {}
for tool_config in enabled_tools:
if tool_config.tool_name is not None:
explicit_by_provider.setdefault(self._provider_id(tool_config), set()).add(tool_config.tool_name)
expanded: list[AgentSoulDifyToolConfig] = []
for tool_config in enabled_tools:
if tool_config.tool_name is not None:
expanded.append(tool_config)
continue
provider_id = self._provider_id(tool_config)
try:
tool_names = self._provider_tools_lister(tenant_id=tenant_id, provider_id=provider_id)
except ToolProviderNotFoundError as exc:
raise WorkflowAgentPluginToolsBuildError(
"agent_tool_declaration_not_found",
f"Dify Plugin Tool provider {provider_id!r} declaration not found: {exc}",
) from exc
if not tool_names:
raise WorkflowAgentPluginToolsBuildError(
"agent_tool_declaration_not_found",
f"Dify Plugin Tool provider {provider_id!r} declares no tools.",
)
already_explicit = explicit_by_provider.get(provider_id, set())
for tool_name in tool_names:
if tool_name in already_explicit:
continue
expanded.append(tool_config.model_copy(update={"tool_name": tool_name, "runtime_parameters": {}}))
return expanded
def _fetch_tool_runtime(
self,
*,
tenant_id: str,
app_id: str,
user_id: str | None,
agent_tool: AgentToolEntity,
invoke_from: InvokeFrom,
tool_config: AgentSoulDifyToolConfig,
) -> Tool:
"""Resolve the API-side ``Tool`` runtime, mapping fetch errors to
Inspector-friendly error codes so callers can render distinct UX for
"tool definition gone" vs "credential failed".
"""
try:
return self._tool_runtime_provider.get_agent_tool_runtime(
tenant_id=tenant_id,
app_id=app_id,
agent_tool=agent_tool,
user_id=user_id,
invoke_from=invoke_from,
variable_pool=None,
allow_file_parameters=True,
use_default_for_missing_form_parameters=True,
)
except ToolProviderNotFoundError as exc:
raise WorkflowAgentPluginToolsBuildError(
"agent_tool_declaration_not_found",
f"Dify Plugin Tool {tool_config.tool_name!r} declaration not found: {exc}",
) from exc
except ToolProviderCredentialValidationError as exc:
raise WorkflowAgentPluginToolsBuildError(
"agent_tool_credential_invalid",
f"Dify Plugin Tool {tool_config.tool_name!r} credential validation failed: {exc}",
) from exc
except ValueError as exc:
# ToolManager raises bare ValueError when the agent tool's
# ``runtime`` / runtime parameters are missing. Surface it under a
# narrower error code than a generic "declaration not found" so
# frontend can render an actionable hint.
raise WorkflowAgentPluginToolsBuildError(
"agent_tool_config_invalid",
f"Dify Plugin Tool {tool_config.tool_name!r} runtime construction failed: {exc}",
) from exc
@staticmethod
def _to_agent_tool_entity(tool_config: AgentSoulDifyToolConfig) -> AgentToolEntity:
# Provider-level entries are expanded into per-tool clones before this point.
assert tool_config.tool_name is not None
return AgentToolEntity(
provider_type=ToolProviderType.value_of(tool_config.provider_type),
provider_id=WorkflowAgentPluginToolsBuilder._provider_id(tool_config),
tool_name=tool_config.tool_name,
tool_parameters=dict(tool_config.runtime_parameters),
credential_id=tool_config.credential_ref.id if tool_config.credential_ref else None,
)
@staticmethod
def _provider_id(tool_config: AgentSoulDifyToolConfig) -> str:
if tool_config.provider_id:
return tool_config.provider_id
assert tool_config.plugin_id is not None
assert tool_config.provider is not None
return f"{tool_config.plugin_id}/{tool_config.provider}"
@staticmethod
def _exposed_tool_name(tool_config: AgentSoulDifyToolConfig) -> str:
# Stage 3.1 decision: no user rename yet. Keep the model-visible tool
# name aligned with the plugin declaration identity. Provider-level
# entries are expanded into per-tool clones before this point.
assert tool_config.tool_name is not None
return tool_config.tool_name
def _to_backend_tool_config(
self,
tool_config: AgentSoulDifyToolConfig,
tool_runtime: Tool,
exposed_name: str,
) -> DifyPluginToolConfig:
runtime = tool_runtime.runtime
if runtime is None:
raise WorkflowAgentPluginToolsBuildError(
"agent_tool_config_invalid",
f"Dify Plugin Tool {tool_config.tool_name!r} has no runtime.",
)
provider_id = self._provider_id(tool_config)
plugin_id, provider = self._plugin_provider(tool_config, provider_id)
parameters = [
DifyPluginToolParameter.model_validate(parameter.model_dump(mode="json"))
for parameter in tool_runtime.get_merged_runtime_parameters()
]
runtime_parameters = self._runtime_parameters(tool_runtime, parameters)
description = tool_config.description
if description is None and tool_runtime.entity.description is not None:
description = tool_runtime.entity.description.llm
return DifyPluginToolConfig(
plugin_id=plugin_id,
provider=provider,
tool_name=exposed_name,
credential_type=self._credential_type(tool_config, runtime.credentials),
name=exposed_name,
description=description,
credentials=self._normalize_credentials(runtime.credentials, tool_name=exposed_name),
runtime_parameters=runtime_parameters,
parameters=parameters,
parameters_json_schema=tool_runtime.get_llm_parameters_json_schema(),
)
@staticmethod
def _plugin_provider(tool_config: AgentSoulDifyToolConfig, provider_id: str) -> tuple[str, str]:
if tool_config.plugin_id and tool_config.provider:
return tool_config.plugin_id, tool_config.provider
provider_id_entity = ToolProviderID(provider_id)
return provider_id_entity.plugin_id, provider_id_entity.provider_name
@staticmethod
def _credential_type(
tool_config: AgentSoulDifyToolConfig,
credentials: Mapping[str, Any],
) -> DifyPluginToolCredentialType:
if not credentials and tool_config.credential_type == "unauthorized":
return "unauthorized"
return tool_config.credential_type
@staticmethod
def _runtime_parameters(
tool_runtime: Tool,
parameters: list[DifyPluginToolParameter],
) -> dict[str, Any]:
runtime = tool_runtime.runtime
runtime_parameters = dict(runtime.runtime_parameters if runtime is not None else {})
missing = [
parameter.name
for parameter in parameters
if parameter.form is not DifyPluginToolParameterForm.LLM
and parameter.required
and parameter.default is None
and parameter.name not in runtime_parameters
]
if missing:
names = ", ".join(sorted(missing))
raise WorkflowAgentPluginToolsBuildError(
"agent_tool_runtime_parameter_missing",
f"Dify Plugin Tool {tool_runtime.entity.identity.name!r} is missing runtime parameters: {names}.",
)
return runtime_parameters
@staticmethod
def _normalize_credentials(
credentials: Mapping[str, Any],
*,
tool_name: str,
) -> dict[str, DifyPluginCredentialValue]:
"""Forward only scalar credential values to the Agent backend.
``DifyPluginCredentialValue`` is ``str | int | float | bool | None``.
Refusing non-scalar values (lists, dicts, custom objects) is safer than
``str(value)`` — stringifying a nested OAuth token blob produces a
Python ``repr`` that the plugin daemon cannot use, and we'd rather
surface a clear ``agent_tool_credential_shape_invalid`` than send junk.
"""
normalized: dict[str, DifyPluginCredentialValue] = {}
for key, value in credentials.items():
if isinstance(value, str | int | float | bool) or value is None:
normalized[key] = value
continue
raise WorkflowAgentPluginToolsBuildError(
"agent_tool_credential_shape_invalid",
(
f"Dify Plugin Tool {tool_name!r} credential {key!r} has a non-scalar value "
f"({type(value).__name__}); only str/int/float/bool/None are forwarded to the daemon."
),
)
return normalized
@@ -8,9 +8,11 @@ from typing import Any, Literal, Protocol, assert_never, cast
from agenton.compositor import CompositorSessionSnapshot
from dify_agent.agent_stub.protocol import AgentStubFileMapping
from dify_agent.layers.ask_human import DifyAskHumanLayerConfig
from dify_agent.layers.drive import (
DifyDriveLayerConfig,
DifyDriveSkillConfig,
from dify_agent.layers.config import (
DifyConfigFileConfig,
DifyConfigLayerConfig,
DifyConfigSkillConfig,
DifyConfigVersionConfig,
)
from dify_agent.layers.execution_context import (
DifyExecutionContextInvokeFrom,
@@ -55,6 +57,7 @@ from models.agent_config_entities import (
AgentKnowledgeModelConfig,
AgentKnowledgeRetrievalConfig,
AgentSoulConfig,
AgentSoulToolsConfig,
DeclaredArrayItem,
DeclaredOutputChildConfig,
DeclaredOutputConfig,
@@ -76,10 +79,14 @@ from services.agent.prompt_mentions import (
parse_prompt_mentions,
workflow_previous_node_output_refs_from_selectors,
)
from services.agent_drive_service import AgentDriveService, decode_drive_mention_ref
from .dify_tools_builder import (
WorkflowAgentDifyToolLayersBuilder,
WorkflowAgentDifyToolsBuilder,
WorkflowAgentDifyToolsBuildError,
WorkflowAgentToolLayers,
)
from .output_failure_orchestrator import retry_idempotency_key
from .plugin_tools_builder import WorkflowAgentPluginToolsBuilder, WorkflowAgentPluginToolsBuildError
from .runtime_feature_manifest import build_runtime_feature_manifest
_DENIED_PERMISSION_STATUSES = frozenset({"unauthorized", "denied", "forbidden", "invalid", "unavailable"})
@@ -97,6 +104,7 @@ _AGENT_STUB_FILE_TRANSFER_METHODS: Mapping[FileTransferMethod, AgentStubFileTran
FileTransferMethod.DATASOURCE_FILE: "datasource_file",
FileTransferMethod.REMOTE_URL: "remote_url",
}
_CANONICAL_DIFY_FILE_REFERENCE_PATTERN = r"^dify-file-ref:eyJyZWNvcmRfaWQiOi[A-Za-z0-9_-]+={0,2}$"
class WorkflowAgentRuntimeRequestBuildError(ValueError):
@@ -157,11 +165,11 @@ class WorkflowAgentRuntimeRequestBuilder:
*,
credentials_provider: CredentialsProvider,
request_builder: AgentBackendRunRequestBuilder | None = None,
plugin_tools_builder: WorkflowAgentPluginToolsBuilder | None = None,
dify_tools_builder: WorkflowAgentDifyToolLayersBuilder | None = None,
) -> None:
self._credentials_provider = credentials_provider
self._request_builder = request_builder or AgentBackendRunRequestBuilder()
self._plugin_tools_builder = plugin_tools_builder or WorkflowAgentPluginToolsBuilder()
self._dify_tools_builder = dify_tools_builder or WorkflowAgentDifyToolsBuilder()
def build(self, context: WorkflowAgentRuntimeBuildContext) -> WorkflowAgentRuntimeRequest:
agent_soul = AgentSoulConfig.model_validate(context.snapshot.config_snapshot_dict)
@@ -182,42 +190,33 @@ class WorkflowAgentRuntimeRequestBuilder:
user_prompt = workflow_context_prompt or self._WORKFLOW_USER_PROMPT_FALLBACK
credentials = self._credentials_provider.fetch(agent_soul.model.model_provider, agent_soul.model.model)
try:
tools_layer = self._plugin_tools_builder.build(
tool_layers = self._build_tool_layers(
tenant_id=context.dify_context.tenant_id,
app_id=context.dify_context.app_id,
user_id=context.dify_context.user_id,
tools=agent_soul.tools,
# Thread the *real* runtime invocation source through to
# ToolManager so credential quotas, rate limits, and audit
# trails match the actual call site (DEBUGGER for draft test
# run, SERVICE_API / WEB_APP for published run).
invoke_from=context.dify_context.invoke_from,
)
except WorkflowAgentPluginToolsBuildError as error:
except WorkflowAgentDifyToolsBuildError as error:
raise WorkflowAgentRuntimeRequestBuildError(error.error_code, str(error)) from error
if tools_layer is not None or agent_soul.tools.cli_tools:
if tool_layers.plugin_tools is not None or tool_layers.core_tools is not None or agent_soul.tools.cli_tools:
metadata["agent_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 [],
"dify_tool_count": len(tool_layers.exposed_tool_names()),
"dify_tool_names": tool_layers.exposed_tool_names(),
"cli_tool_count": len(agent_soul.tools.cli_tools),
}
drive_config: DifyDriveLayerConfig | None = None
config_layer_config: DifyConfigLayerConfig | None = None
soul_prompt_resolver = build_soul_mention_resolver(agent_soul)
if dify_config.AGENT_DRIVE_MANIFEST_ENABLED:
drive_config, drive_warnings = build_drive_layer_config(
config_layer_config, config_warnings = build_config_layer_config(
agent_soul,
tenant_id=context.dify_context.tenant_id,
agent_id=context.agent.id,
)
append_runtime_warnings(metadata, drive_warnings)
soul_prompt_resolver = build_drive_aware_soul_mention_resolver(
agent_soul,
tenant_id=context.dify_context.tenant_id,
agent_id=context.agent.id,
config_version_id=context.snapshot.id,
config_version_kind="snapshot",
)
append_runtime_warnings(metadata, config_warnings)
soul_prompt_resolver = build_config_aware_soul_mention_resolver(agent_soul)
soul_prompt = expand_prompt_mentions(agent_soul.prompt.system_prompt, soul_prompt_resolver).strip()
knowledge_config = build_knowledge_layer_config(agent_soul)
@@ -251,6 +250,7 @@ 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,
agent_config_version_kind="snapshot",
agent_mode=self._agent_backend_agent_mode(context.dify_context.invoke_from),
invoke_from=cast(DifyExecutionContextInvokeFrom, context.dify_context.invoke_from.value),
),
@@ -258,9 +258,10 @@ class WorkflowAgentRuntimeRequestBuilder:
workflow_node_job_prompt=workflow_job_prompt,
user_prompt=user_prompt,
output=self._build_output_config(node_job.declared_outputs),
tools=tools_layer,
tools=tool_layers.plugin_tools,
core_tools=tool_layers.core_tools,
knowledge=knowledge_config,
drive_config=drive_config,
config_layer_config=config_layer_config,
ask_human_config=build_ask_human_layer_config(agent_soul),
include_shell=dify_config.AGENT_SHELL_ENABLED,
shell_config=build_shell_layer_config(agent_soul),
@@ -279,6 +280,26 @@ class WorkflowAgentRuntimeRequestBuilder:
metadata=metadata,
)
def _build_tool_layers(
self,
*,
tenant_id: str,
app_id: str,
user_id: str | None,
tools: AgentSoulToolsConfig,
invoke_from: InvokeFrom,
) -> WorkflowAgentToolLayers:
# Production workflow runs intentionally keep existing plugin configs on
# the direct `dify.plugin.tools` route. This builder emits plugin tools
# directly and non-plugin Dify tools through `dify.core.tools`.
return self._dify_tools_builder.build_layers(
tenant_id=tenant_id,
app_id=app_id,
user_id=user_id,
tools=tools,
invoke_from=invoke_from,
)
@staticmethod
def _agent_backend_agent_mode(invoke_from: InvokeFrom) -> Literal["workflow_run", "single_step"]:
if invoke_from in {InvokeFrom.DEBUGGER, InvokeFrom.VALIDATION}:
@@ -494,7 +515,10 @@ class WorkflowAgentRuntimeRequestBuilder:
schema: dict[str, Any] = {"type": "object", "properties": properties}
if required:
schema["required"] = required
return AgentBackendOutputConfig(json_schema=schema)
return AgentBackendOutputConfig(
json_schema=schema,
description=WorkflowAgentRuntimeRequestBuilder._build_output_description(effective_outputs),
)
@staticmethod
def effective_declared_outputs(
@@ -551,48 +575,107 @@ class WorkflowAgentRuntimeRequestBuilder:
array_schema["items"]["description"] = array_item.description
return array_schema
case DeclaredOutputType.FILE:
return {
"oneOf": [
{
"type": "object",
"additionalProperties": False,
"properties": {
"transfer_method": {"const": FileTransferMethod.LOCAL_FILE.value},
"reference": {"type": "string"},
},
"required": ["transfer_method", "reference"],
},
{
"type": "object",
"additionalProperties": False,
"properties": {
"transfer_method": {"const": FileTransferMethod.TOOL_FILE.value},
"reference": {"type": "string"},
},
"required": ["transfer_method", "reference"],
},
{
"type": "object",
"additionalProperties": False,
"properties": {
"transfer_method": {"const": FileTransferMethod.DATASOURCE_FILE.value},
"reference": {"type": "string"},
},
"required": ["transfer_method", "reference"],
},
{
"type": "object",
"additionalProperties": False,
"properties": {
"transfer_method": {"const": FileTransferMethod.REMOTE_URL.value},
"url": {"type": "string"},
},
"required": ["transfer_method", "url"],
},
],
}
return WorkflowAgentRuntimeRequestBuilder._agent_stub_output_file_mapping_schema()
assert_never(output_type)
@staticmethod
def _agent_stub_output_file_mapping_schema() -> dict[str, Any]:
"""JSON Schema for Agent-produced output file mappings.
``AgentStubFileMapping.model_json_schema()`` cannot express the model's
``after`` validator: only ``remote_url`` may carry ``url``; every other
method must carry a canonical ``reference``. The structured-output model
needs that relationship in the schema, otherwise it may emit
``{"transfer_method": "local_file", "url": "..."}``, which passes the
broad generated schema but fails API-side output type checking.
For files produced inside an Agent run, the supported persisted shape is
narrower than every downloadable mapping: the sandbox must upload the
local artifact via ``dify-agent file upload <path>``, which returns a
``tool_file`` mapping. ``local_file`` and ``datasource_file`` are valid
for existing file references in workflow context, not for newly produced
Agent output files.
"""
return {
"title": "AgentStubFileMapping",
"description": (
"Agent output file mapping. Use `tool_file` with `reference` for files uploaded by "
"`dify-agent file upload <path>`; use `remote_url` only for files already reachable by URL."
),
"anyOf": [
{
"type": "object",
"additionalProperties": False,
"required": ["transfer_method", "reference"],
"properties": {
"transfer_method": {
"type": "string",
"enum": ["tool_file"],
},
"reference": {
"type": "string",
"minLength": 1,
"pattern": _CANONICAL_DIFY_FILE_REFERENCE_PATTERN,
"description": (
"Canonical Dify file reference returned by `dify-agent file upload <path>`. "
"Never use a local path, filename, URL, or synthesized dify-file-ref here."
),
},
},
},
{
"type": "object",
"additionalProperties": False,
"required": ["transfer_method", "url"],
"properties": {
"transfer_method": {
"type": "string",
"enum": ["remote_url"],
},
"url": {
"type": "string",
"minLength": 1,
"description": "Remote URL for a file that is already publicly reachable.",
},
},
},
],
}
@staticmethod
def _build_output_description(declared_outputs: Sequence[DeclaredOutputConfig]) -> str | None:
file_output_lines: list[str] = []
for output in declared_outputs:
if output.type == DeclaredOutputType.FILE:
file_output_lines.append(
f"- `{output.name}`: create the file in the sandbox, run `dify-agent file upload <path>`, "
f"and set `final_output.{output.name}` to the returned AgentStubFileMapping JSON object. "
"Do not call `final_output` before the upload command succeeds. Do not use the local path, "
"filename, URL, or a synthesized/base64-encoded value as the `reference`."
)
elif (
output.type == DeclaredOutputType.ARRAY
and output.array_item is not None
and output.array_item.type == DeclaredOutputType.FILE
):
file_output_lines.append(
f"- `{output.name}`: for every produced file, run `dify-agent file upload <path>` and set "
f"`final_output.{output.name}` to an array of the returned AgentStubFileMapping JSON objects. "
"Do not call `final_output` before all upload commands succeed. Do not use local paths, filenames, "
"URLs, or synthesized/base64-encoded values as `reference` values."
)
if not file_output_lines:
return None
return "\n".join(
[
"When filling file outputs, do not return a local filesystem path directly.",
"Upload each sandbox-local file through the Agent Stub CLI first. Copy the JSON printed by "
"`dify-agent file upload <path>` verbatim into the final output; never invent the `reference` value.",
*file_output_lines,
]
)
@staticmethod
def _apply_child_properties(schema: dict[str, Any], children: Sequence[DeclaredOutputChildConfig]) -> None:
if not children:
@@ -753,19 +836,12 @@ def append_runtime_warnings(metadata: dict[str, Any], warnings: list[dict[str, s
existing.extend(warnings)
def build_drive_aware_soul_mention_resolver(
agent_soul: AgentSoulConfig,
*,
tenant_id: str,
agent_id: str,
):
"""Resolve skill/file mentions against the agent drive and everything else via Agent Soul."""
def build_config_aware_soul_mention_resolver(agent_soul: AgentSoulConfig):
"""Resolve config skill/file mentions and delegate the rest to Agent Soul."""
base_resolver = build_soul_mention_resolver(agent_soul)
drive_service = AgentDriveService()
skill_catalog = drive_service.list_skills(tenant_id=tenant_id, agent_id=agent_id)
skill_names_by_key = {skill["skill_md_key"]: skill["name"] for skill in skill_catalog}
drive_keys = {item["key"] for item in drive_service.manifest(tenant_id=tenant_id, agent_id=agent_id)}
skill_names = {item.name for item in agent_soul.config_skills}
file_names = {item.name for item in agent_soul.config_files}
def _resolve(mention: object) -> str | None:
if not hasattr(mention, "kind") or not hasattr(mention, "ref_id"):
@@ -774,88 +850,97 @@ def build_drive_aware_soul_mention_resolver(
ref_id = cast(str, mention.ref_id)
label = cast(str | None, getattr(mention, "label", None))
if kind == MentionKind.SKILL:
decoded_key = decode_drive_mention_ref(ref_id)
return skill_names_by_key.get(decoded_key) or label or decoded_key
return ref_id if ref_id in skill_names else label or ref_id
if kind == MentionKind.FILE:
decoded_key = decode_drive_mention_ref(ref_id)
if decoded_key in drive_keys:
return decoded_key.rsplit("/", 1)[-1]
return label or decoded_key
return ref_id if ref_id in file_names else label or ref_id
return base_resolver(cast(Any, mention))
return _resolve
def build_drive_layer_config(
def build_config_layer_config(
agent_soul: AgentSoulConfig,
*,
tenant_id: str,
agent_id: str | None,
) -> tuple[DifyDriveLayerConfig | None, list[dict[str, str]]]:
"""Derive drive runtime catalog + prompt-mentioned eager-pull keys from the drive."""
agent_id: str | None = None,
config_version_id: str | None = None,
config_version_kind: Literal["snapshot", "draft", "build_draft"] = "snapshot",
) -> tuple[DifyConfigLayerConfig | None, list[dict[str, str]]]:
"""Derive prompt-mentioned eager-pull names from Agent Soul."""
mentioned_drive_refs = [
decode_drive_mention_ref(mention.ref_id)
for mention in parse_prompt_mentions(agent_soul.prompt.system_prompt)
if mention.kind in {MentionKind.SKILL, MentionKind.FILE}
]
ordered_mentions = list(dict.fromkeys(ref for ref in mentioned_drive_refs if ref))
if not agent_id:
if not ordered_mentions:
return None, []
return None, [
{
"section": "agent_soul.prompt.system_prompt",
"code": "drive_ref_dangling",
"message": "drive mentions are configured but the run has no bound agent to address a drive by.",
}
]
ordered_mentions = list(
dict.fromkeys(
mention.ref_id
for mention in parse_prompt_mentions(agent_soul.prompt.system_prompt)
if mention.kind in {MentionKind.SKILL, MentionKind.FILE} and mention.ref_id
)
)
if (
not agent_soul.config_skills
and not agent_soul.config_files
and not agent_soul.config_note
and not ordered_mentions
):
return None, []
drive_service = AgentDriveService()
skills_catalog = drive_service.list_skills(tenant_id=tenant_id, agent_id=agent_id)
manifest_items = drive_service.manifest(tenant_id=tenant_id, agent_id=agent_id)
manifest_by_key = {item["key"]: item for item in manifest_items}
skill_keys = {skill["skill_md_key"] for skill in skills_catalog}
skill_names = {skill.name for skill in agent_soul.config_skills}
file_names = {file_ref.name for file_ref in agent_soul.config_files}
warnings: list[dict[str, str]] = []
mentioned_skill_keys: list[str] = []
mentioned_file_keys: list[str] = []
for drive_key in ordered_mentions:
if drive_key in skill_keys:
mentioned_skill_keys.append(drive_key)
mentioned_skill_names: list[str] = []
mentioned_file_names: list[str] = []
for name in ordered_mentions:
if name in skill_names:
mentioned_skill_names.append(name)
continue
if drive_key in manifest_by_key:
mentioned_file_keys.append(drive_key)
if name in file_names:
mentioned_file_names.append(name)
continue
warnings.append(
{
"section": "agent_soul.prompt.system_prompt",
"code": "mention_target_missing",
"message": f"drive mention '{drive_key}' has no matching drive entry.",
"message": f"config mention '{name}' has no matching config asset.",
}
)
skills = [
DifyDriveSkillConfig(
path=skill["path"],
name=skill["name"],
description=skill["description"],
skill_md_key=skill["skill_md_key"],
archive_key=skill["archive_key"],
)
for skill in skills_catalog
]
return (
DifyDriveLayerConfig(
drive_ref=f"agent-{agent_id}",
skills=skills,
mentioned_skill_keys=mentioned_skill_keys,
mentioned_file_keys=mentioned_file_keys,
DifyConfigLayerConfig(
agent_id=agent_id,
config_version=DifyConfigVersionConfig(
id=config_version_id,
kind=config_version_kind,
writable=config_version_kind == "build_draft",
),
skills=[
DifyConfigSkillConfig(
name=skill.name,
description=skill.description,
size=skill.size,
mime_type=skill.mime_type,
)
for skill in agent_soul.config_skills
],
files=[
DifyConfigFileConfig(
name=file_ref.name,
size=file_ref.size,
mime_type=file_ref.mime_type,
)
for file_ref in agent_soul.config_files
],
env_keys=_agent_soul_config_env_keys(agent_soul),
note=agent_soul.config_note,
mentioned_skill_names=mentioned_skill_names,
mentioned_file_names=mentioned_file_names,
),
warnings,
)
def _agent_soul_config_env_keys(agent_soul: AgentSoulConfig) -> list[str]:
keys = [item.key or item.name or item.env_name or item.variable for item in agent_soul.env.variables]
return [key for key in keys if key]
def _cli_tool_enabled(item: object) -> bool:
"""A CLI tool is bootstrapped unless explicitly disabled (default is enabled)."""
data = _plain_mapping(item)