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

Signed-off-by: yyh <yuanyouhuilyz@gmail.com>
Co-authored-by: 盐粒 Yanli <mail@yanli.one>
Co-authored-by: Joel <iamjoel007@gmail.com>
Co-authored-by: zyssyz123 <916125788@qq.com>
Co-authored-by: autofix-ci[bot] <114827586+autofix-ci[bot]@users.noreply.github.com>
Co-authored-by: 盐粒 Yanli <yanli@dify.ai>
Co-authored-by: 林玮 (Jade Lin) <linw1995@icloud.com>
This commit is contained in:
yyh
2026-07-05 08:09:38 +00:00
committed by GitHub
co-authored by 盐粒 Yanli Joel zyssyz123 autofix-ci[bot] 盐粒 Yanli 林玮
parent f8d47616c1
commit 1b3e8e9943
303 changed files with 11391 additions and 1912 deletions
@@ -507,6 +507,7 @@ class AgentAppGenerator(MessageBasedAppGenerator):
),
event_adapter=AgentBackendRunEventAdapter(),
session_store=AgentAppRuntimeSessionStore(),
text_delta_debounce_seconds=dify_config.AGENT_APP_TEXT_DELTA_DEBOUNCE_SECONDS,
)
def _run_input_guards(
+102 -9
View File
@@ -16,6 +16,8 @@ from __future__ import annotations
import json
import logging
import time
from collections.abc import Mapping
from decimal import Decimal
from typing import Any, Literal
@@ -74,12 +76,23 @@ def _prompt_messages_from_query(user_query: str | None) -> list[PromptMessage]:
return [UserPromptMessage(content=user_query)]
def _llm_usage_from_agent_backend(usage: Mapping[str, Any] | None) -> LLMUsage | None:
if usage is None:
return None
try:
return LLMUsage.from_metadata(usage)
except (TypeError, ValueError):
logger.warning("Failed to parse Agent backend usage metadata: %s", usage, exc_info=True)
return None
def publish_text_answer(
*,
queue_manager: AppQueueManager,
model_name: str,
answer: str,
user_query: str | None = None,
usage: LLMUsage | None = None,
) -> None:
"""Publish a complete assistant answer as one chunk + message-end.
@@ -99,6 +112,7 @@ def publish_text_answer(
model_name=model_name,
answer=answer,
user_query=user_query,
usage=usage,
)
@@ -127,6 +141,7 @@ def publish_message_end(
model_name: str,
answer: str,
user_query: str | None = None,
usage: LLMUsage | None = None,
) -> None:
"""Publish the terminal assistant result without emitting another delta."""
prompt_messages = _prompt_messages_from_query(user_query)
@@ -136,13 +151,46 @@ def publish_message_end(
model=model_name,
prompt_messages=prompt_messages,
message=AssistantPromptMessage(content=answer),
usage=LLMUsage.empty_usage(),
usage=usage or LLMUsage.empty_usage(),
),
),
PublishFrom.APPLICATION_MANAGER,
)
class _TextDeltaDebouncer:
"""Batch assistant text deltas on stream-event boundaries for final SSE output."""
def __init__(self, *, debounce_seconds: float) -> None:
self._debounce_seconds = debounce_seconds
self._parts: list[str] = []
self._first_pending_at: float | None = None
def push(self, delta: str) -> str | None:
if not delta:
return None
if self._debounce_seconds <= 0:
return delta
now = time.monotonic()
if not self._parts:
self._first_pending_at = now
self._parts.append(delta)
if self._first_pending_at is not None and now - self._first_pending_at >= self._debounce_seconds:
return self.flush()
return None
def flush(self) -> str | None:
if not self._parts:
return None
text = "".join(self._parts)
self._parts = []
self._first_pending_at = None
return text
class _AgentProcessRecorder:
"""Persist Agent v2 thinking/tool process events through the legacy thought model."""
@@ -430,11 +478,13 @@ class AgentAppRunner:
agent_backend_client: AgentBackendRunClient,
event_adapter: AgentBackendRunEventAdapter,
session_store: AgentAppRuntimeSessionStore,
text_delta_debounce_seconds: float,
) -> None:
self._request_builder = request_builder
self._agent_backend_client = agent_backend_client
self._event_adapter = event_adapter
self._session_store = session_store
self._text_delta_debounce_seconds = text_delta_debounce_seconds
def run(
self,
@@ -512,6 +562,7 @@ class AgentAppRunner:
answer=answer,
query=query,
streamed_answer=streamed_answer,
usage=_llm_usage_from_agent_backend(terminal.usage),
)
self._save_session(
scope=scope,
@@ -724,19 +775,39 @@ class AgentAppRunner:
model_name: str,
query: str | None,
):
"""Consume backend events while preserving raw recorder granularity.
Process events are recorded immediately for observability. Only the
final assistant text deltas sent through the EasyUI queue are debounced,
with flushes happening on later stream events or terminal boundaries.
"""
terminal = None
streamed_answer_parts: list[str] = []
text_delta_debouncer = _TextDeltaDebouncer(debounce_seconds=self._text_delta_debounce_seconds)
process_recorder = _AgentProcessRecorder(
dify_context=dify_context,
message_id=message_id,
queue_manager=queue_manager,
)
def flush_pending_text() -> None:
pending_text = text_delta_debouncer.flush()
if pending_text:
publish_text_delta(
queue_manager=queue_manager,
model_name=model_name,
delta=pending_text,
user_query=query,
)
for public_event in self._agent_backend_client.stream_events(run_id):
if queue_manager.is_stopped():
flush_pending_text()
self._cancel_run(run_id)
raise GenerateTaskStoppedError()
for internal_event in self._event_adapter.adapt(public_event):
if queue_manager.is_stopped():
flush_pending_text()
self._cancel_run(run_id)
raise GenerateTaskStoppedError()
if internal_event.type in (
@@ -758,18 +829,22 @@ class AgentAppRunner:
text_delta = self._extract_stream_text_delta(internal_event)
if text_delta:
streamed_answer_parts.append(text_delta)
publish_text_delta(
queue_manager=queue_manager,
model_name=model_name,
delta=text_delta,
user_query=query,
)
debounced_delta = text_delta_debouncer.push(text_delta)
if debounced_delta:
publish_text_delta(
queue_manager=queue_manager,
model_name=model_name,
delta=debounced_delta,
user_query=query,
)
continue
continue
flush_pending_text()
terminal = internal_event
break
if terminal is not None:
break
flush_pending_text()
return terminal, "".join(streamed_answer_parts)
def _cancel_run(self, run_id: str) -> None:
@@ -793,10 +868,20 @@ class AgentAppRunner:
answer: str,
query: str | None,
streamed_answer: str,
usage: LLMUsage | None,
) -> None:
"""Finish a successful streamed turn without duplicating the final text."""
if not answer and streamed_answer:
answer = streamed_answer
if not streamed_answer:
self._publish_answer(queue_manager=queue_manager, model_name=model_name, answer=answer, query=query)
publish_text_answer(
queue_manager=queue_manager,
model_name=model_name,
answer=answer,
user_query=query,
usage=usage,
)
return
if answer.startswith(streamed_answer):
@@ -812,7 +897,13 @@ class AgentAppRunner:
"using terminal output for message persistence."
)
publish_message_end(queue_manager=queue_manager, model_name=model_name, answer=answer, user_query=query)
publish_message_end(
queue_manager=queue_manager,
model_name=model_name,
answer=answer,
user_query=query,
usage=usage,
)
def _save_session(
self,
@@ -852,6 +943,8 @@ class AgentAppRunner:
configured the value is a JSON object, which we serialize so the chat
message always has a string body.
"""
if output is None:
return ""
if isinstance(output, str):
return output
if isinstance(output, dict):
@@ -46,7 +46,7 @@ from core.workflow.nodes.agent_v2.runtime_request_builder import (
)
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
from services.agent.prompt_mentions import expand_prompt_mentions
class AgentAppRuntimeRequestBuildError(ValueError):
@@ -124,17 +124,14 @@ class AgentAppRuntimeRequestBuilder:
"cli_tool_count": len(agent_soul.tools.cli_tools),
}
config_layer_config = None
soul_prompt_resolver = build_soul_mention_resolver(agent_soul)
if dify_config.AGENT_DRIVE_MANIFEST_ENABLED:
config_layer_config, config_warnings = build_config_layer_config(
agent_soul,
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)
config_layer_config, config_warnings = build_config_layer_config(
agent_soul,
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(
+3 -1
View File
@@ -1,5 +1,6 @@
import logging
from collections.abc import Sequence
from urllib.parse import urlencode
import httpx
from yarl import URL
@@ -16,7 +17,8 @@ MARKETPLACE_TIMEOUT = 30
def get_plugin_pkg_url(plugin_unique_identifier: str) -> str:
return str((marketplace_api_url / "api/v1/plugins/download").with_query(unique_identifier=plugin_unique_identifier))
query = urlencode({"unique_identifier": plugin_unique_identifier})
return f"{marketplace_api_url / 'api/v1/plugins/download'}?{query}"
def download_plugin_pkg(plugin_unique_identifier: str):
+1
View File
@@ -230,6 +230,7 @@ class RequestRequestUploadFile(BaseModel):
filename: str
mimetype: str
conversation_id: str | None = None
class RequestDownloadFileMapping(BaseModel):
+4 -1
View File
@@ -902,7 +902,10 @@ class PluginService:
tenant_id,
plugin_unique_identifiers,
PluginInstallationSource.Package,
[{}],
[
{"plugin_unique_identifier": plugin_unique_identifier}
for plugin_unique_identifier in plugin_unique_identifiers
],
)
PluginService.invalidate_plugin_model_providers_cache(tenant_id)
return result
+36 -1
View File
@@ -1159,6 +1159,39 @@ class DatasetRetrieval:
all_documents.extend(documents)
def _run_retriever_thread(
self,
*,
flask_app: Flask,
dataset_id: str,
query: str | None,
top_k: int,
all_documents: list[Document],
document_ids_filter: list[str] | None,
metadata_condition: MetadataFilteringCondition | None,
attachment_ids: list[str] | None,
cancel_event: threading.Event | None,
thread_exceptions: list[Exception] | None,
) -> None:
try:
with session_factory.create_session() as session:
self._retriever(
flask_app=flask_app,
session=session,
dataset_id=dataset_id,
query=query or "",
top_k=top_k,
all_documents=all_documents,
document_ids_filter=document_ids_filter,
metadata_condition=metadata_condition,
attachment_ids=attachment_ids,
)
except Exception as e:
if cancel_event:
cancel_event.set()
if thread_exceptions is not None:
thread_exceptions.append(e)
def to_dataset_retriever_tool(
self,
session: Session,
@@ -1797,7 +1830,7 @@ class DatasetRetrieval:
else:
continue
retrieval_thread = threading.Thread(
target=self._retriever,
target=self._run_retriever_thread,
kwargs={
"flask_app": flask_app,
"dataset_id": dataset.id,
@@ -1807,6 +1840,8 @@ class DatasetRetrieval:
"document_ids_filter": document_ids_filter,
"metadata_condition": metadata_condition,
"attachment_ids": [attachment_id] if attachment_id else None,
"cancel_event": cancel_event,
"thread_exceptions": thread_exceptions,
},
)
threads.append(retrieval_thread)
+24 -13
View File
@@ -64,34 +64,45 @@ def verify_tool_file_signature(file_id: str, timestamp: str, nonce: str, sign: s
return current_time - int(timestamp) <= dify_config.FILES_ACCESS_TIMEOUT
def get_signed_file_url_for_plugin(filename: str, mimetype: str, tenant_id: str, user_id: str) -> str:
def get_signed_file_url_for_plugin(
filename: str, mimetype: str, tenant_id: str, user_id: str, conversation_id: str | None = None
) -> str:
"""Build the signed upload URL used by the plugin-facing file upload endpoint."""
base_url = dify_config.INTERNAL_FILES_URL or dify_config.FILES_URL
upload_url = f"{base_url}/files/upload/for-plugin"
timestamp = str(int(time.time()))
nonce = os.urandom(16).hex()
data_to_sign = f"upload|{filename}|{mimetype}|{tenant_id}|{user_id}|{timestamp}|{nonce}"
data_to_sign = f"upload|{filename}|{mimetype}|{tenant_id}|{user_id}|{conversation_id or ''}|{timestamp}|{nonce}"
sign = hmac.new(_secret_key(), data_to_sign.encode(), hashlib.sha256).digest()
encoded_sign = base64.urlsafe_b64encode(sign).decode()
query = urllib.parse.urlencode(
{
"timestamp": timestamp,
"nonce": nonce,
"sign": encoded_sign,
"user_id": user_id,
"tenant_id": tenant_id,
}
)
query_params = {
"timestamp": timestamp,
"nonce": nonce,
"sign": encoded_sign,
"user_id": user_id,
"tenant_id": tenant_id,
}
if conversation_id:
query_params["conversation_id"] = conversation_id
query = urllib.parse.urlencode(query_params)
return f"{upload_url}?{query}"
def verify_plugin_file_signature(
*, filename: str, mimetype: str, tenant_id: str, user_id: str, timestamp: str, nonce: str, sign: str
*,
filename: str,
mimetype: str,
tenant_id: str,
user_id: str,
conversation_id: str | None = None,
timestamp: str,
nonce: str,
sign: str,
) -> bool:
"""Verify the signature used by the plugin-facing file upload endpoint."""
data_to_sign = f"upload|{filename}|{mimetype}|{tenant_id}|{user_id}|{timestamp}|{nonce}"
data_to_sign = f"upload|{filename}|{mimetype}|{tenant_id}|{user_id}|{conversation_id or ''}|{timestamp}|{nonce}"
recalculated_sign = hmac.new(_secret_key(), data_to_sign.encode(), hashlib.sha256).digest()
recalculated_encoded_sign = base64.urlsafe_b64encode(recalculated_sign).decode()
@@ -11,6 +11,7 @@ from dify_agent.layers.dify_plugin import (
DifyPluginToolCredentialType,
DifyPluginToolParameter,
DifyPluginToolParameterForm,
DifyPluginToolParameterType,
DifyPluginToolsLayerConfig,
)
from sqlalchemy import select
@@ -351,7 +352,9 @@ class WorkflowAgentDifyToolsBuilder:
@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:
if provider_type is ToolProviderType.PLUGIN or (
provider_type is ToolProviderType.BUILT_IN and _is_plugin_provider_id(tool_config.provider_id)
):
return "plugin"
if provider_type in {
ToolProviderType.BUILT_IN,
@@ -404,7 +407,7 @@ class WorkflowAgentDifyToolsBuilder:
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(),
parameters_json_schema=self._plugin_parameters_json_schema(tool_runtime, parameters),
)
def _to_core_backend_tool_config(
@@ -456,6 +459,41 @@ class WorkflowAgentDifyToolsBuilder:
description = tool_runtime.entity.description.llm
return description
@staticmethod
def _plugin_parameters_json_schema(
tool_runtime: Tool,
parameters: list[DifyPluginToolParameter],
) -> dict[str, Any]:
schema = tool_runtime.get_llm_parameters_json_schema()
properties = schema.setdefault("properties", {})
required = schema.setdefault("required", [])
if not isinstance(properties, dict) or not isinstance(required, list):
raise WorkflowAgentDifyToolsBuildError(
"agent_tool_declaration_invalid",
f"Dify Plugin Tool {tool_runtime.entity.identity.name!r} has invalid parameter schema.",
)
for parameter in parameters:
if parameter.form is not DifyPluginToolParameterForm.LLM:
continue
if parameter.type is DifyPluginToolParameterType.FILE:
properties[parameter.name] = _plugin_file_input_schema(parameter.llm_description or "")
elif parameter.type in {
DifyPluginToolParameterType.FILES,
DifyPluginToolParameterType.SYSTEM_FILES,
}:
properties[parameter.name] = {
"type": "array",
"items": _plugin_file_input_schema(parameter.llm_description or ""),
"description": parameter.llm_description or "",
}
else:
continue
if parameter.required and parameter.name not in required:
required.append(parameter.name)
return schema
@staticmethod
def _runtime_parameters(
tool_runtime: Tool,
@@ -498,3 +536,44 @@ class WorkflowAgentDifyToolsBuilder:
),
)
return normalized
def _is_plugin_provider_id(provider_id: str | None) -> bool:
if not provider_id:
return False
parts = provider_id.split("/")
return len(parts) == 3 and all(parts)
def _plugin_file_input_schema(description: str) -> dict[str, Any]:
return {
"description": description,
"anyOf": [
{
"type": "string",
"minLength": 1,
"description": "HTTP(S) URL or sandbox-local file path.",
},
{
"type": "object",
"additionalProperties": False,
"required": ["transfer_method", "url"],
"properties": {
"transfer_method": {"type": "string", "enum": ["remote_url"]},
"url": {"type": "string", "minLength": 1},
},
},
{
"type": "object",
"additionalProperties": False,
"required": ["transfer_method", "reference"],
"properties": {
"transfer_method": {
"type": "string",
"enum": ["local_file", "tool_file", "datasource_file"],
},
"reference": {"type": "string", "minLength": 1},
},
},
],
}
@@ -333,6 +333,8 @@ class WorkflowAgentOutputAdapter:
session_snapshot = None
if isinstance(event, AgentBackendRunSucceededInternalEvent | AgentBackendDeferredToolCallInternalEvent):
session_snapshot = event.session_snapshot
if event.usage is not None:
agent_backend["usage"] = dict(event.usage)
if session_snapshot is not None:
agent_backend["session_snapshot"] = {
"layer_count": len(session_snapshot.layers),
@@ -206,17 +206,14 @@ class WorkflowAgentRuntimeRequestBuilder:
"cli_tool_count": len(agent_soul.tools.cli_tools),
}
config_layer_config: DifyConfigLayerConfig | None = None
soul_prompt_resolver = build_soul_mention_resolver(agent_soul)
if dify_config.AGENT_DRIVE_MANIFEST_ENABLED:
config_layer_config, config_warnings = build_config_layer_config(
agent_soul,
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)
config_layer_config, config_warnings = build_config_layer_config(
agent_soul,
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)
@@ -717,7 +714,7 @@ def build_shell_layer_config(agent_soul: AgentSoulConfig) -> DifyShellLayerConfi
for tool in (_shell_cli_tool(item) for item in agent_soul.tools.cli_tools if _cli_tool_enabled(item))
if tool is not None
],
env=[env for env in (_shell_env_var(item) for item in agent_soul.env.variables) if env is not None],
env=_shell_env_vars(agent_soul.env.variables, agent_soul.env.secret_refs),
secret_refs=[
secret for secret in (_shell_secret_ref(item) for item in agent_soul.env.secret_refs) if secret is not None
],
@@ -864,8 +861,13 @@ def build_config_layer_config(
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."""
) -> tuple[DifyConfigLayerConfig, list[dict[str, str]]]:
"""Build the always-present Agent config layer from Agent Soul state.
The ``dify.config`` layer must exist for every Agent v2 runtime request so
the backend can expose the config CLI/help surface even when the current
Agent Soul has no config assets, note, or prompt mentions.
"""
ordered_mentions = list(
dict.fromkeys(
@@ -874,14 +876,6 @@ def build_config_layer_config(
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, []
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]] = []
@@ -968,11 +962,7 @@ def _shell_cli_tool(item: object) -> DifyShellCliToolConfig | None:
if not commands and not isinstance(name, str):
return None
tool_env = data.get("env") if isinstance(data.get("env"), Mapping) else {}
env = [
env_var
for env_var in (_shell_env_var(item) for item in _env_entries(tool_env, "variables"))
if env_var is not None
]
env = _shell_env_vars(_env_entries(tool_env, "variables"), _env_entries(tool_env, "secret_refs"))
secret_refs = [
secret_ref
for secret_ref in (_shell_secret_ref(item) for item in _env_entries(tool_env, "secret_refs"))
@@ -995,6 +985,12 @@ def _env_entries(env: object, key: str) -> list[object]:
return entries
def _shell_env_vars(variables: Sequence[object], secret_refs: Sequence[object]) -> list[DifyShellEnvVarConfig]:
env_vars = [_shell_env_var(item) for item in variables]
secret_env_vars = [_shell_env_var(item) for item in secret_refs if _has_secret_value(item)]
return [env for env in [*env_vars, *secret_env_vars] if env is not None]
def _shell_env_var(item: object) -> DifyShellEnvVarConfig | None:
data = _plain_mapping(item)
name = _name_from_mapping(data)
@@ -1011,13 +1007,15 @@ def _shell_secret_ref(item: object) -> DifyShellSecretRefConfig | None:
name = _name_from_mapping(data)
if name is None:
return None
ref = (
data.get("ref")
or data.get("value")
or data.get("id")
or data.get("credential_id")
or data.get("provider_credential_id")
)
# Inline Composer values are passed as env vars because the agent-backend
# secret ref schema only accepts short backend-managed reference IDs.
if _has_secret_value(item):
return None
ref = data.get("ref") or data.get("credential_id") or data.get("provider_credential_id")
if ref is None:
ref = data.get("id")
if ref is None:
return None
return DifyShellSecretRefConfig(name=name, ref=str(ref) if ref is not None else None)
@@ -1029,6 +1027,12 @@ def _plain_mapping(item: object) -> dict[str, Any]:
return {}
def _has_secret_value(item: object) -> bool:
data = _plain_mapping(item)
value = data.get("value")
return isinstance(value, str) and bool(value)
def _name_from_mapping(item: Mapping[str, Any]) -> str | None:
for key in ("name", "key", "env_name", "variable"):
value = item.get(key)