diff --git a/api/.env.example b/api/.env.example index 2adde29d334..104eded9513 100644 --- a/api/.env.example +++ b/api/.env.example @@ -569,8 +569,8 @@ WORKFLOW_GENERATOR_NODE_BUILDER_MAX_WORKERS=6 GRAPH_ENGINE_MIN_WORKERS=3 # Maximum number of workers per GraphEngine instance (default: 10) GRAPH_ENGINE_MAX_WORKERS=10 -# Queue depth threshold that triggers worker scale up (default: 3) -GRAPH_ENGINE_SCALE_UP_THRESHOLD=3 +# Pending task threshold that triggers worker scale up (default: 0) +GRAPH_ENGINE_SCALE_UP_THRESHOLD=0 # Seconds of idle time before scaling down workers (default: 5.0) GRAPH_ENGINE_SCALE_DOWN_IDLE_TIME=5.0 diff --git a/api/configs/feature/__init__.py b/api/configs/feature/__init__.py index 12971b9b8ca..e979a9486c3 100644 --- a/api/configs/feature/__init__.py +++ b/api/configs/feature/__init__.py @@ -856,9 +856,9 @@ class WorkflowConfig(BaseSettings): default=10, ) - GRAPH_ENGINE_SCALE_UP_THRESHOLD: PositiveInt = Field( - description="Queue depth threshold that triggers worker scale up", - default=3, + GRAPH_ENGINE_SCALE_UP_THRESHOLD: NonNegativeInt = Field( + description="Pending task threshold that triggers worker scale up", + default=0, ) GRAPH_ENGINE_SCALE_DOWN_IDLE_TIME: float = Field( diff --git a/api/core/workflow/node_factory.py b/api/core/workflow/node_factory.py index 4f8f0553ab4..f98e221e61c 100644 --- a/api/core/workflow/node_factory.py +++ b/api/core/workflow/node_factory.py @@ -361,7 +361,6 @@ class DifyNodeFactory(NodeFactory): self._agent_runtime_support = AgentRuntimeSupport() self._agent_message_transformer = AgentMessageTransformer() - @override def with_runtime_state(self, graph_runtime_state: "GraphRuntimeState") -> "DifyNodeFactory": return DifyNodeFactory( graph_init_params=self.graph_init_params, diff --git a/api/core/workflow/workflow_entry.py b/api/core/workflow/workflow_entry.py index 493890e814d..dd504bb27da 100644 --- a/api/core/workflow/workflow_entry.py +++ b/api/core/workflow/workflow_entry.py @@ -38,6 +38,7 @@ from graphon.graph_engine.layers import DebugLoggingLayer, ExecutionLimitsLayer from graphon.graph_events import GraphEngineEvent, GraphNodeEventBase, GraphRunFailedEvent from graphon.nodes import BuiltinNodeTypes from graphon.nodes.base.node import Node +from graphon.nodes.container_effects import ContainerAwaitRequest from graphon.runtime import GraphRuntimeState, VariablePool from graphon.variable_loader import DUMMY_VARIABLE_LOADER, VariableLoader, load_into_variable_pool from models.workflow import Workflow @@ -199,7 +200,7 @@ class WorkflowEntry: user_inputs: Mapping[str, Any], variable_pool: VariablePool, variable_loader: VariableLoader = DUMMY_VARIABLE_LOADER, - ) -> tuple[Node, Generator[GraphNodeEventBase, None, None]]: + ) -> tuple[Node, Generator[GraphNodeEventBase | ContainerAwaitRequest, None, None]]: """ Single step run workflow node :param workflow: Workflow instance @@ -347,7 +348,7 @@ class WorkflowEntry: @classmethod def run_free_node( cls, node_data: dict[str, Any], node_id: str, tenant_id: str, user_id: str, user_inputs: dict[str, Any] - ) -> tuple[Node, Generator[GraphNodeEventBase, None, None]]: + ) -> tuple[Node, Generator[GraphNodeEventBase | ContainerAwaitRequest, None, None]]: """ Run free node @@ -541,7 +542,7 @@ class WorkflowEntry: variable_pool.add([variable_node_id] + variable_key_list, input_value) @staticmethod - def _traced_node_run(node: Node) -> Generator[GraphNodeEventBase, None, None]: + def _traced_node_run(node: Node) -> Generator[GraphNodeEventBase | ContainerAwaitRequest, None, None]: """ Wraps a node's run method with OpenTelemetry tracing and returns a generator. """ diff --git a/api/pyproject.toml b/api/pyproject.toml index 88bc598c79a..c7910433be9 100644 --- a/api/pyproject.toml +++ b/api/pyproject.toml @@ -45,7 +45,7 @@ dependencies = [ "zstandard==0.25.0", # Emerging: newer and fast-moving, use compatible pins "fastopenapi[flask]==0.7.0", - "graphon==0.6.0", + "graphon==0.7.0", "httpx-sse==0.4.3", "json-repair==0.60.1", ] @@ -67,7 +67,7 @@ exclude = ["providers/vdb/__pycache__", "providers/trace/__pycache__"] [tool.uv.sources] dify-agent = { path = "../dify-agent", editable = true } flask-restx = { git = "https://github.com/asukaminato0721/flask-restx", rev = "27758e26f8f740d7525d5039c51a9e524b6e2b68" } -graphon = { git = "https://github.com/langgenius/graphon", rev = "6d85e98df87a74589303dcb297c7866196359d27" } +graphon = { git = "https://github.com/langgenius/graphon", rev = "d48c36fb02d8aa0d31dc6a9140a27c04a370600f" } dify-vdb-alibabacloud-mysql = { workspace = true } dify-vdb-analyticdb = { workspace = true } dify-vdb-baidu = { workspace = true } diff --git a/api/services/rag_pipeline/rag_pipeline.py b/api/services/rag_pipeline/rag_pipeline.py index 6d7f3a01af4..6df864064f5 100644 --- a/api/services/rag_pipeline/rag_pipeline.py +++ b/api/services/rag_pipeline/rag_pipeline.py @@ -49,6 +49,7 @@ from graphon.errors import WorkflowNodeRunFailedError from graphon.graph_events import GraphNodeEventBase, NodeRunFailedEvent, NodeRunSucceededEvent from graphon.node_events import NodeRunResult from graphon.nodes.base.node import Node +from graphon.nodes.container_effects import ContainerAwaitRequest from graphon.nodes.http_request import HTTP_REQUEST_CONFIG_FILTER_KEY, build_http_request_config from graphon.runtime import VariablePool from graphon.variables.variables import Variable, VariableBase @@ -909,7 +910,10 @@ class RagPipelineService: def _handle_node_run_result( self, - getter: Callable[[], tuple[Node, Generator[GraphNodeEventBase, None, None]]], + getter: Callable[ + [], + tuple[Node, Generator[GraphNodeEventBase | ContainerAwaitRequest, None, None]], + ], start_at: float, tenant_id: str, node_id: str, diff --git a/api/services/workflow_service.py b/api/services/workflow_service.py index b29affe696a..c53ee872fed 100644 --- a/api/services/workflow_service.py +++ b/api/services/workflow_service.py @@ -62,6 +62,7 @@ from graphon.graph_events import GraphNodeEventBase, NodeRunFailedEvent, NodeRun from graphon.node_events import NodeRunResult from graphon.nodes import BuiltinNodeTypes from graphon.nodes.base.node import Node +from graphon.nodes.container_effects import ContainerAwaitRequest from graphon.nodes.http_request import HTTP_REQUEST_CONFIG_FILTER_KEY, build_http_request_config from graphon.nodes.start.entities import StartNodeData from graphon.runtime import VariablePool @@ -1447,7 +1448,10 @@ class WorkflowService: def _handle_single_step_result( self, - invoke_node_fn: Callable[[], tuple[Node, Generator[GraphNodeEventBase, None, None]]], + invoke_node_fn: Callable[ + [], + tuple[Node, Generator[GraphNodeEventBase | ContainerAwaitRequest, None, None]], + ], start_at: float, node_id: str, ) -> WorkflowNodeExecution: @@ -1483,7 +1487,11 @@ class WorkflowService: return node_execution def _execute_node_safely( - self, invoke_node_fn: Callable[[], tuple[Node, Generator[GraphNodeEventBase, None, None]]] + self, + invoke_node_fn: Callable[ + [], + tuple[Node, Generator[GraphNodeEventBase | ContainerAwaitRequest, None, None]], + ], ) -> tuple[Node, NodeRunResult | None, bool, str | None]: """ Execute node safely and handle errors according to error strategy. diff --git a/api/tests/unit_tests/configs/test_dify_config.py b/api/tests/unit_tests/configs/test_dify_config.py index 8807de47ff6..e02a3828256 100644 --- a/api/tests/unit_tests/configs/test_dify_config.py +++ b/api/tests/unit_tests/configs/test_dify_config.py @@ -78,6 +78,7 @@ def test_dify_config(monkeypatch: pytest.MonkeyPatch): assert config.AGENT_SHELL_ENABLED is True assert config.SENTRY_TRACES_SAMPLE_RATE == 1.0 assert config.TEMPLATE_TRANSFORM_MAX_LENGTH == 400_000 + assert config.GRAPH_ENGINE_SCALE_UP_THRESHOLD == 0 # annotated field with custom configured value assert config.HTTP_REQUEST_MAX_READ_TIMEOUT == 300 diff --git a/api/tests/unit_tests/core/variables/test_segment.py b/api/tests/unit_tests/core/variables/test_segment.py index bf593b52895..f65e5bbde75 100644 --- a/api/tests/unit_tests/core/variables/test_segment.py +++ b/api/tests/unit_tests/core/variables/test_segment.py @@ -28,6 +28,7 @@ from graphon.variables.segments import ( StringSegment, get_segment_discriminator, ) +from graphon.variables.template_resolution import convert_template from graphon.variables.types import SegmentType from graphon.variables.utils import ( dumps_with_segments, diff --git a/api/tests/unit_tests/core/workflow/test_human_input_adapter.py b/api/tests/unit_tests/core/workflow/test_human_input_adapter.py index 7a6328ffb4b..30352738d0b 100644 --- a/api/tests/unit_tests/core/workflow/test_human_input_adapter.py +++ b/api/tests/unit_tests/core/workflow/test_human_input_adapter.py @@ -18,12 +18,12 @@ from core.workflow.human_input_adapter import ( ) from graphon.enums import BuiltinNodeTypes from graphon.nodes.base.variable_template_parser import VariableTemplateParser +from graphon.runtime import VariablePool def test_email_delivery_config_helpers_render_and_sanitize_text() -> None: - variable_pool = SimpleNamespace( - convert_template=lambda body: SimpleNamespace(text=body.replace("{{#node.value#}}", "42")) - ) + variable_pool = VariablePool() + variable_pool.add(["node", "value"], "42") rendered = EmailDeliveryConfig.render_body_template( body="Open {{#url#}} and use {{#node.value#}}", diff --git a/api/uv.lock b/api/uv.lock index 95ff4f8e136..ca65a0b4e38 100644 --- a/api/uv.lock +++ b/api/uv.lock @@ -1638,7 +1638,7 @@ requires-dist = [ { name = "gmpy2", specifier = ">=2.3.0,<3.0.0" }, { name = "google-api-python-client", specifier = ">=2.198.0,<3.0.0" }, { name = "google-cloud-aiplatform", specifier = ">=1.160.0,<2.0.0" }, - { name = "graphon", git = "https://github.com/langgenius/graphon?rev=6d85e98df87a74589303dcb297c7866196359d27" }, + { name = "graphon", git = "https://github.com/langgenius/graphon?rev=d48c36fb02d8aa0d31dc6a9140a27c04a370600f" }, { name = "gunicorn", specifier = ">=26.0.0,<27.0.0" }, { name = "httpx", extras = ["socks"], specifier = "==0.28.1" }, { name = "httpx-sse", specifier = "==0.4.3" }, @@ -2991,8 +2991,8 @@ httpx = [ [[package]] name = "graphon" -version = "0.6.0" -source = { git = "https://github.com/langgenius/graphon?rev=6d85e98df87a74589303dcb297c7866196359d27#6d85e98df87a74589303dcb297c7866196359d27" } +version = "0.7.0" +source = { git = "https://github.com/langgenius/graphon?rev=d48c36fb02d8aa0d31dc6a9140a27c04a370600f#d48c36fb02d8aa0d31dc6a9140a27c04a370600f" } dependencies = [ { name = "charset-normalizer" }, { name = "httpx" }, diff --git a/docker/envs/core-services/shared.env.example b/docker/envs/core-services/shared.env.example index bf2c32d5b1d..b39d50b2ad5 100644 --- a/docker/envs/core-services/shared.env.example +++ b/docker/envs/core-services/shared.env.example @@ -192,7 +192,7 @@ WORKFLOW_GENERATION_TIMEOUT_MS=180000 WORKFLOW_FILE_UPLOAD_LIMIT=10 GRAPH_ENGINE_MIN_WORKERS=3 GRAPH_ENGINE_MAX_WORKERS=10 -GRAPH_ENGINE_SCALE_UP_THRESHOLD=3 +GRAPH_ENGINE_SCALE_UP_THRESHOLD=0 GRAPH_ENGINE_SCALE_DOWN_IDLE_TIME=5.0 ALIYUN_SLS_ACCESS_KEY_ID= ALIYUN_SLS_ACCESS_KEY_SECRET=