Compare commits

..
233 changed files with 3742 additions and 6801 deletions
+6 -6
View File
@@ -29,13 +29,13 @@ jobs:
steps:
- name: Checkout code
uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1
uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0
with:
fetch-depth: 0
persist-credentials: false
- name: Setup UV and Python
uses: astral-sh/setup-uv@c771a70e6277c0a99b617c7a806ffedaca235ff9 # v9.0.0
uses: astral-sh/setup-uv@11f9893b081a58869d3b5fccaea48c9e9e46f990 # v8.3.2
with:
enable-cache: true
python-version: ${{ matrix.python-version }}
@@ -88,13 +88,13 @@ jobs:
steps:
- name: Checkout code
uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1
uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0
with:
fetch-depth: 0
persist-credentials: false
- name: Setup UV and Python
uses: astral-sh/setup-uv@c771a70e6277c0a99b617c7a806ffedaca235ff9 # v9.0.0
uses: astral-sh/setup-uv@11f9893b081a58869d3b5fccaea48c9e9e46f990 # v8.3.2
with:
enable-cache: true
python-version: ${{ matrix.python-version }}
@@ -139,13 +139,13 @@ jobs:
steps:
- name: Checkout code
uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1
uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0
with:
fetch-depth: 0
persist-credentials: false
- name: Setup UV and Python
uses: astral-sh/setup-uv@c771a70e6277c0a99b617c7a806ffedaca235ff9 # v9.0.0
uses: astral-sh/setup-uv@11f9893b081a58869d3b5fccaea48c9e9e46f990 # v8.3.2
with:
enable-cache: true
python-version: "3.12"
+3 -3
View File
@@ -20,7 +20,7 @@ jobs:
run: echo "autofix.ci updates pull request branches, not merge group refs."
- if: github.event_name != 'merge_group'
uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1
uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0
- name: Check Docker Compose inputs
if: github.event_name != 'merge_group'
@@ -84,12 +84,12 @@ jobs:
dify-agent/pyproject.toml
dify-agent/uv.lock
- if: github.event_name != 'merge_group'
uses: actions/setup-python@5fda3b95a4ea91299a34e894583c3862153e4b97 # v7.0.0
uses: actions/setup-python@ece7cb06caefa5fff74198d8649806c4678c61a1 # v6.3.0
with:
python-version: "3.11"
- if: github.event_name != 'merge_group'
uses: astral-sh/setup-uv@c771a70e6277c0a99b617c7a806ffedaca235ff9 # v9.0.0
uses: astral-sh/setup-uv@11f9893b081a58869d3b5fccaea48c9e9e46f990 # v8.3.2
- name: Generate Docker Compose
if: github.event_name != 'merge_group' && steps.docker-compose-changes.outputs.any_changed == 'true'
+2 -2
View File
@@ -97,7 +97,7 @@ jobs:
echo "PLATFORM_PAIR=${platform//\//-}" >> $GITHUB_ENV
- name: Login to Docker Hub
uses: docker/login-action@abd2ef45e78c5afb21d64d4ca52ee8550d9572c7 # v4.5.1
uses: docker/login-action@af1e73f918a031802d376d3c8bbc3fe56130a9b0 # v4.4.0
with:
username: ${{ env.DOCKERHUB_USER }}
password: ${{ env.DOCKERHUB_TOKEN }}
@@ -199,7 +199,7 @@ jobs:
merge-multiple: true
- name: Login to Docker Hub
uses: docker/login-action@abd2ef45e78c5afb21d64d4ca52ee8550d9572c7 # v4.5.1
uses: docker/login-action@af1e73f918a031802d376d3c8bbc3fe56130a9b0 # v4.4.0
with:
username: ${{ env.DOCKERHUB_USER }}
password: ${{ env.DOCKERHUB_TOKEN }}
+6 -6
View File
@@ -79,7 +79,7 @@ jobs:
ws2_app_id: ${{ steps.out.outputs.DIFY_E2E_WS2_APP_ID }}
steps:
- uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v4
- uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v4
with:
ref: ${{ inputs.cli_ref || github.ref }}
persist-credentials: false
@@ -123,7 +123,7 @@ jobs:
shell: bash
steps:
- uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v4
- uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v4
with:
ref: ${{ inputs.cli_ref || github.ref }}
persist-credentials: false
@@ -170,7 +170,7 @@ jobs:
shell: bash
steps:
- uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v4
- uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v4
with:
ref: ${{ inputs.cli_ref || github.ref }}
persist-credentials: false
@@ -233,7 +233,7 @@ jobs:
shell: bash
steps:
- uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v4
- uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v4
with:
ref: ${{ inputs.cli_ref || github.ref }}
persist-credentials: false
@@ -295,7 +295,7 @@ jobs:
shell: bash
steps:
- uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v4
- uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v4
with:
ref: ${{ inputs.cli_ref || github.ref }}
persist-credentials: false
@@ -351,7 +351,7 @@ jobs:
shell: bash
steps:
- uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v4
- uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v4
with:
ref: ${{ inputs.cli_ref || github.ref }}
persist-credentials: false
+1 -1
View File
@@ -23,7 +23,7 @@ jobs:
working-directory: ./cli
steps:
- name: Checkout
uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1
uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0
with:
persist-credentials: false
fetch-depth: 0
+2 -2
View File
@@ -35,7 +35,7 @@ jobs:
dify_tag: ${{ steps.resolve.outputs.dify_tag }}
steps:
- name: Checkout
uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1
uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0
with:
persist-credentials: false
@@ -98,7 +98,7 @@ jobs:
DIFY_TAG: ${{ needs.validate.outputs.dify_tag }}
steps:
- name: Checkout
uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1
uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0
with:
persist-credentials: false
fetch-depth: 1
+1 -1
View File
@@ -24,7 +24,7 @@ jobs:
shell: bash
steps:
- name: Checkout cli ref
uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1
uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0
with:
ref: ${{ inputs.cli_ref || github.ref }}
persist-credentials: false
+1 -1
View File
@@ -30,7 +30,7 @@ jobs:
steps:
- name: Checkout code
uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1
uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0
with:
persist-credentials: false
+4 -4
View File
@@ -13,13 +13,13 @@ jobs:
steps:
- name: Checkout code
uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1
uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0
with:
fetch-depth: 0
persist-credentials: false
- name: Setup UV and Python
uses: astral-sh/setup-uv@c771a70e6277c0a99b617c7a806ffedaca235ff9 # v9.0.0
uses: astral-sh/setup-uv@11f9893b081a58869d3b5fccaea48c9e9e46f990 # v8.3.2
with:
enable-cache: true
python-version: "3.12"
@@ -63,13 +63,13 @@ jobs:
steps:
- name: Checkout code
uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1
uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0
with:
fetch-depth: 0
persist-credentials: false
- name: Setup UV and Python
uses: astral-sh/setup-uv@c771a70e6277c0a99b617c7a806ffedaca235ff9 # v9.0.0
uses: astral-sh/setup-uv@11f9893b081a58869d3b5fccaea48c9e9e46f990 # v8.3.2
with:
enable-cache: true
python-version: "3.12"
+1 -1
View File
@@ -24,7 +24,7 @@ jobs:
name: Require cherry-pick provenance
runs-on: depot-ubuntu-24.04
steps:
- uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1
- uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0
with:
fetch-depth: 0
+1 -1
View File
@@ -9,6 +9,6 @@ jobs:
pull-requests: write
runs-on: depot-ubuntu-24.04
steps:
- uses: actions/labeler@bf12e9b00b37c5c0ca2b87b79b2daf7891dbda13 # v7.0.0
- uses: actions/labeler@b8dd2d9be0f68b860e7dae5dae7d772984eacd6d # v6.2.0
with:
sync-labels: true
+1 -1
View File
@@ -47,7 +47,7 @@ jobs:
migration-changed: ${{ steps.changes.outputs.migration }}
sandbox-runtime-changed: ${{ steps.changes.outputs.sandbox-runtime }}
steps:
- uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1
- uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0
- uses: dorny/paths-filter@7b450fff21473bca461d4b92ce414b9d0420d706 # v4.0.2
id: changes
with:
+1 -1
View File
@@ -18,7 +18,7 @@ jobs:
outputs:
external-e2e-changed: ${{ steps.changes.outputs.external_e2e }}
steps:
- uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1
- uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0
- uses: dorny/paths-filter@7b450fff21473bca461d4b92ce414b9d0420d706 # v4.0.2
id: changes
with:
+2 -2
View File
@@ -17,12 +17,12 @@ jobs:
pull-requests: write
steps:
- name: Checkout PR branch
uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1
uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0
with:
fetch-depth: 0
- name: Setup Python & UV
uses: astral-sh/setup-uv@c771a70e6277c0a99b617c7a806ffedaca235ff9 # v9.0.0
uses: astral-sh/setup-uv@11f9893b081a58869d3b5fccaea48c9e9e46f990 # v8.3.2
with:
enable-cache: true
@@ -21,10 +21,10 @@ jobs:
if: ${{ github.event.workflow_run.conclusion == 'success' && github.event.workflow_run.pull_requests[0].head.repo.full_name != github.repository }}
steps:
- name: Checkout default branch (trusted code)
uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1
uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0
- name: Setup Python & UV
uses: astral-sh/setup-uv@c771a70e6277c0a99b617c7a806ffedaca235ff9 # v9.0.0
uses: astral-sh/setup-uv@11f9893b081a58869d3b5fccaea48c9e9e46f990 # v8.3.2
with:
enable-cache: true
+2 -2
View File
@@ -17,12 +17,12 @@ jobs:
pull-requests: write
steps:
- name: Checkout PR branch
uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1
uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0
with:
fetch-depth: 0
- name: Setup Python & UV
uses: astral-sh/setup-uv@c771a70e6277c0a99b617c7a806ffedaca235ff9 # v9.0.0
uses: astral-sh/setup-uv@11f9893b081a58869d3b5fccaea48c9e9e46f990 # v8.3.2
with:
enable-cache: true
+3 -3
View File
@@ -21,7 +21,7 @@ jobs:
steps:
- name: Checkout code
uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1
uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0
with:
fetch-depth: 0
persist-credentials: false
@@ -45,7 +45,7 @@ jobs:
steps:
- name: Checkout code
uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1
uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0
with:
fetch-depth: 0
persist-credentials: false
@@ -72,7 +72,7 @@ jobs:
steps:
- name: Checkout code
uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1
uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0
with:
fetch-depth: 0
persist-credentials: false
+5 -5
View File
@@ -23,7 +23,7 @@ jobs:
steps:
- name: Checkout code
uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1
uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0
with:
persist-credentials: false
fetch-depth: 0
@@ -45,7 +45,7 @@ jobs:
- name: Setup UV and Python
if: steps.changed-files.outputs.any_changed == 'true'
uses: astral-sh/setup-uv@c771a70e6277c0a99b617c7a806ffedaca235ff9 # v9.0.0
uses: astral-sh/setup-uv@11f9893b081a58869d3b5fccaea48c9e9e46f990 # v8.3.2
with:
enable-cache: false
python-version: "3.12"
@@ -93,7 +93,7 @@ jobs:
steps:
- name: Checkout code
uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1
uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0
with:
persist-credentials: false
@@ -144,7 +144,7 @@ jobs:
steps:
- name: Checkout code
uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1
uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0
with:
persist-credentials: false
@@ -186,7 +186,7 @@ jobs:
steps:
- name: Checkout code
uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1
uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0
with:
fetch-depth: 0
persist-credentials: false
+1 -1
View File
@@ -24,7 +24,7 @@ jobs:
working-directory: sdks/nodejs-client
steps:
- uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1
- uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0
with:
persist-credentials: false
+2 -2
View File
@@ -40,7 +40,7 @@ jobs:
steps:
- name: Checkout repository
uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1
uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0
with:
fetch-depth: 0
token: ${{ secrets.GITHUB_TOKEN }}
@@ -158,7 +158,7 @@ jobs:
- name: Run Claude Code for Translation Sync
if: steps.context.outputs.CHANGED_FILES != ''
uses: anthropics/claude-code-action@be7b93b1907a4abad570368f3c74b6fe3807510b # v1.0.183
uses: anthropics/claude-code-action@af0559ee4f514d1ef21826982bed13f7edc3c35e # v1.0.178
with:
anthropic_api_key: ${{ secrets.ANTHROPIC_API_KEY }}
github_token: ${{ secrets.GITHUB_TOKEN }}
+1 -1
View File
@@ -21,7 +21,7 @@ jobs:
steps:
- name: Checkout repository
uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1
uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0
with:
fetch-depth: 0
+2 -2
View File
@@ -24,7 +24,7 @@ jobs:
steps:
- name: Checkout code
uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1
uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0
with:
persist-credentials: false
@@ -36,7 +36,7 @@ jobs:
remove_tool_cache: true
- name: Setup UV and Python
uses: astral-sh/setup-uv@c771a70e6277c0a99b617c7a806ffedaca235ff9 # v9.0.0
uses: astral-sh/setup-uv@11f9893b081a58869d3b5fccaea48c9e9e46f990 # v8.3.2
with:
enable-cache: true
python-version: ${{ matrix.python-version }}
+2 -2
View File
@@ -21,7 +21,7 @@ jobs:
steps:
- name: Checkout code
uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1
uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0
with:
persist-credentials: false
@@ -33,7 +33,7 @@ jobs:
remove_tool_cache: true
- name: Setup UV and Python
uses: astral-sh/setup-uv@c771a70e6277c0a99b617c7a806ffedaca235ff9 # v9.0.0
uses: astral-sh/setup-uv@11f9893b081a58869d3b5fccaea48c9e9e46f990 # v8.3.2
with:
enable-cache: true
python-version: ${{ matrix.python-version }}
+2 -2
View File
@@ -26,7 +26,7 @@ jobs:
steps:
- name: Checkout code
uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1
uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0
with:
persist-credentials: false
@@ -34,7 +34,7 @@ jobs:
uses: ./.github/actions/setup-web
- name: Setup UV and Python
uses: astral-sh/setup-uv@c771a70e6277c0a99b617c7a806ffedaca235ff9 # v9.0.0
uses: astral-sh/setup-uv@11f9893b081a58869d3b5fccaea48c9e9e46f990 # v8.3.2
with:
enable-cache: true
python-version: "3.12"
+4 -4
View File
@@ -29,7 +29,7 @@ jobs:
steps:
- name: Checkout code
uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1
uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0
with:
persist-credentials: false
@@ -62,7 +62,7 @@ jobs:
steps:
- name: Checkout code
uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1
uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0
with:
persist-credentials: false
@@ -100,7 +100,7 @@ jobs:
steps:
- name: Checkout code
uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1
uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0
with:
persist-credentials: false
@@ -132,7 +132,7 @@ jobs:
steps:
- name: Checkout code
uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1
uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0
with:
persist-credentials: false
+1 -3
View File
@@ -666,7 +666,6 @@ PLUGIN_REMOTE_INSTALL_PORT=5003
PLUGIN_REMOTE_INSTALL_HOST=localhost
PLUGIN_MAX_PACKAGE_SIZE=15728640
PLUGIN_MODEL_SCHEMA_CACHE_TTL=3600
PLUGIN_MODEL_PROVIDERS_CACHE_ENABLED=true
PLUGIN_MODEL_PROVIDERS_CACHE_TTL=86400
# Comma-separated marketplace plugin IDs whose latest versions are installed for newly registered users.
# Example: langgenius/openai,langgenius/gemini
@@ -678,8 +677,6 @@ INNER_API_KEY_FOR_PLUGIN=QaHbTe77CtuXmsfyhR7+vRjI/+XbV1AaFy691iy+kGDv2Jvy0/eAh8Y
# Dify Agent backend
AGENT_BACKEND_BASE_URL=http://localhost:5050
# Bearer token sent to the Agent backend /runs API. Must match DIFY_AGENT_API_TOKEN on the server side.
AGENT_BACKEND_API_TOKEN=dify-agent-run-token-for-dev-only
AGENT_BACKEND_STREAM_READ_TIMEOUT_SECONDS=30
AGENT_BACKEND_STREAM_MAX_RECONNECTS=3
AGENT_BACKEND_RUN_TIMEOUT_SECONDS=1200
@@ -732,6 +729,7 @@ OTEL_MAX_EXPORT_BATCH_SIZE=512
OTEL_METRIC_EXPORT_INTERVAL=60000
OTEL_BATCH_EXPORT_TIMEOUT=10000
OTEL_METRIC_EXPORT_TIMEOUT=30000
# Prevent Clickjacking
ALLOW_EMBED=false
+30 -59
View File
@@ -1,13 +1,10 @@
import logging
import time
from collections.abc import Callable
from typing import NamedTuple
import socketio
from flask import request
from opentelemetry.trace import get_current_span
from opentelemetry.trace.span import INVALID_SPAN_ID, INVALID_TRACE_ID
from werkzeug.exceptions import Forbidden, HTTPException, ServiceUnavailable
from configs import dify_config
from contexts.wrapper import RecyclableContextVar
@@ -45,53 +42,6 @@ _CONSOLE_EXEMPT_PREFIXES = (
"/console/api/activate/check",
)
_WEBAPP_EXEMPT_PREFIXES = ("/api/system-features",)
_INVALID_LICENSE_STATUSES = (LicenseStatus.INACTIVE, LicenseStatus.EXPIRED, LicenseStatus.LOST)
def _session_surface_error(license_status: LicenseStatus | None) -> HTTPException:
if license_status is None:
return UnauthorizedAndForceLogout("Unable to verify enterprise license. Please contact your administrator.")
return UnauthorizedAndForceLogout(f"Enterprise license is {license_status}. Please contact your administrator.")
def _bearer_surface_error(license_status: LicenseStatus | None) -> HTTPException:
"""Token-authed: forcing a logout is meaningless and license state must not leak."""
return Forbidden(description="license_required")
def _retryable_surface_error(license_status: LicenseStatus | None) -> HTTPException:
"""Webhook senders retry on 5xx but treat 4xx as permanent, disabling the subscription."""
return ServiceUnavailable(description="license_required")
class _LicenseGatedSurface(NamedTuple):
prefix: str
exempt_prefixes: tuple[str, ...]
build_error: Callable[[LicenseStatus | None], HTTPException]
# /files (plugin-daemon data plane), /inner/api (enterprise control plane) and /health
# stay ungated: blocking them breaks workflow execution or license recovery itself.
_LICENSE_GATED_SURFACES = (
_LicenseGatedSurface("/console/api/", _CONSOLE_EXEMPT_PREFIXES, _session_surface_error),
_LicenseGatedSurface("/api/", _WEBAPP_EXEMPT_PREFIXES, _session_surface_error),
_LicenseGatedSurface("/v1", (), _bearer_surface_error),
_LicenseGatedSurface("/mcp", (), _bearer_surface_error),
_LicenseGatedSurface("/triggers", (), _retryable_surface_error),
)
def _match_license_gated_surface(path: str) -> _LicenseGatedSurface | None:
for surface in _LICENSE_GATED_SURFACES:
if not path.startswith(surface.prefix):
continue
if any(path.startswith(exempt) for exempt in surface.exempt_prefixes):
return None
return surface
return None
# ----------------------------
# Application Factory Function
@@ -112,17 +62,38 @@ def create_flask_app_with_configs() -> DifyApp:
init_request_context()
RecyclableContextVar.increment_thread_recycles()
# Enterprise license validation for API endpoints (both console and webapp)
# When license expires, block all API access except bootstrap endpoints needed
# for the frontend to load the license expiration page without infinite reloads.
if dify_config.ENTERPRISE_ENABLED:
surface = _match_license_gated_surface(request.path)
if surface is not None:
try:
license_status = EnterpriseService.get_cached_license_status()
except Exception:
logger.exception("Failed to check enterprise license status")
license_status = None
is_console_api = request.path.startswith("/console/api/")
is_webapp_api = request.path.startswith("/api/")
if license_status is None or license_status in _INVALID_LICENSE_STATUSES:
raise surface.build_error(license_status)
if is_console_api or is_webapp_api:
if is_console_api:
is_exempt = any(request.path.startswith(p) for p in _CONSOLE_EXEMPT_PREFIXES)
else: # webapp API
is_exempt = request.path.startswith("/api/system-features")
if not is_exempt:
try:
# Check license status (cached — see EnterpriseService for TTL details)
license_status = EnterpriseService.get_cached_license_status()
if license_status in (LicenseStatus.INACTIVE, LicenseStatus.EXPIRED, LicenseStatus.LOST):
raise UnauthorizedAndForceLogout(
f"Enterprise license is {license_status}. Please contact your administrator."
)
if license_status is None:
raise UnauthorizedAndForceLogout(
"Unable to verify enterprise license. Please contact your administrator."
)
except UnauthorizedAndForceLogout:
raise
except Exception:
logger.exception("Failed to check enterprise license status")
raise UnauthorizedAndForceLogout(
"Unable to verify enterprise license. Please contact your administrator."
)
# add after request hook for injecting trace headers from OpenTelemetry span context
# Only adds headers when OTEL is enabled and has valid context
+1 -5
View File
@@ -11,7 +11,6 @@ from clients.agent_backend.fake_client import FakeAgentBackendRunClient, FakeAge
def create_agent_backend_run_client(
*,
base_url: str | None = None,
api_token: str | None = None,
use_fake: bool = False,
fake_scenario: str | FakeAgentBackendScenario = FakeAgentBackendScenario.SUCCESS,
stream_read_timeout_seconds: float = 30,
@@ -23,11 +22,8 @@ def create_agent_backend_run_client(
return FakeAgentBackendRunClient(scenario=FakeAgentBackendScenario(fake_scenario))
if base_url is None:
raise ValueError("base_url is required when creating a real Agent backend client")
headers: dict[str, str] = {}
if api_token:
headers["Authorization"] = f"Bearer {api_token}"
return DifyAgentBackendRunClient(
Client(base_url=base_url, stream_timeout=stream_read_timeout_seconds, headers=headers),
Client(base_url=base_url, stream_timeout=stream_read_timeout_seconds),
stream_max_reconnects=stream_max_reconnects,
stream_timeout_seconds=stream_run_timeout_seconds,
)
@@ -12,11 +12,6 @@ class AgentBackendConfig(BaseSettings):
default=None,
)
AGENT_BACKEND_API_TOKEN: str | None = Field(
description="Bearer token for authenticating with the Agent backend /runs API.",
default=None,
)
AGENT_BACKEND_USE_FAKE: bool = Field(
description="Use the deterministic in-process fake Agent backend client.",
default=False,
-42
View File
@@ -266,12 +266,6 @@ class PluginConfig(BaseSettings):
default=60 * 60,
)
PLUGIN_MODEL_PROVIDERS_CACHE_ENABLED: bool = Field(
description="Whether tenant plugin model providers are cached in Redis. Disable when plugins are installed "
"by a system other than this one, which cannot invalidate the cache when a tenant's plugins change.",
default=True,
)
PLUGIN_MODEL_PROVIDERS_CACHE_TTL: PositiveInt = Field(
description="TTL in seconds for caching tenant plugin model providers in Redis",
default=60 * 60 * 24,
@@ -822,41 +816,6 @@ class UpdateConfig(BaseSettings):
)
class CommunityTelemetryConfig(BaseSettings):
"""
Configuration for anonymous self-hosted community telemetry.
"""
DISABLE_TELEMETRY: bool = Field(
description="Disable anonymous community telemetry",
default=False,
)
DO_NOT_TRACK: bool = Field(
description="Respect the standard do-not-track opt-out signal for telemetry",
default=False,
)
TELEMETRY_ENDPOINT: str = Field(
description="Endpoint for anonymous community telemetry events",
default="https://otel.dify.ai/v1/events",
)
TELEMETRY_FALLBACK_ENDPOINT: str = Field(
description="Fallback endpoint for anonymous community telemetry events",
default="https://otel.dify.cn/v1/events",
)
TELEMETRY_TIMEOUT_SECONDS: PositiveInt = Field(
description="HTTP timeout in seconds for anonymous community telemetry requests",
default=3,
)
TELEMETRY_HEARTBEAT_INTERVAL_MINUTES: PositiveInt = Field(
description="Celery beat interval in minutes for checking whether heartbeat telemetry is due",
default=30,
)
CI: bool = Field(
description="Whether the process is running in CI; telemetry is skipped when true",
default=False,
)
class WorkflowVariableTruncationConfig(BaseSettings):
WORKFLOW_VARIABLE_TRUNCATION_MAX_SIZE: PositiveInt = Field(
# 1000 KiB
@@ -1640,7 +1599,6 @@ class FeatureConfig(
TenantIsolatedTaskQueueConfig,
ToolConfig,
UpdateConfig,
CommunityTelemetryConfig,
WorkflowConfig,
WorkflowNodeExecutionConfig,
WorkspaceConfig,
+5
View File
@@ -257,6 +257,7 @@ class AgentAppDetailWithSite(GenericAppDetailWithSite):
debug_conversation_has_messages: bool = False
debug_conversation_message_count: int = 0
role: str | None = None
active_config_is_published: bool = False
class AgentDebugConversationRefreshResponse(BaseModel):
@@ -409,6 +410,10 @@ def _serialize_agent_app_detail(
payload["debug_conversation_has_messages"] = message_count > 0
payload["debug_conversation_message_count"] = message_count
payload["role"] = agent.role or ""
payload["active_config_is_published"] = roster_service.active_config_is_published(
tenant_id=app_model.tenant_id,
agent=agent,
)
return payload
-77
View File
@@ -58,7 +58,6 @@ from services.app_service import (
AppResponseView,
AppService,
CreateAppParams,
RecentAppMode,
StarredAppListParams,
)
from services.enterprise import rbac_service as enterprise_rbac_service
@@ -140,10 +139,6 @@ class AppListBaseQuery(BaseModel):
raise ValueError("Invalid UUID format in creator_ids.") from exc
class RecentAppListQuery(BaseModel):
limit: int = Field(default=8, ge=1, le=8, description="Number of recently modified apps to return (1-8)")
class AppListQuery(AppListBaseQuery):
pass
@@ -416,33 +411,6 @@ class AppPartial(AppResponseModel):
return to_timestamp(value)
class RecentAppResponse(ResponseModel):
id: str
name: str
icon_type: IconType | None = None
icon: str | None = None
icon_background: str | None = None
mode: RecentAppMode
author_name: str | None = None
updated_at: int
permission_keys: list[str] = Field(default_factory=list)
maintainer: str | None = None
@computed_field(return_type=str | None) # type: ignore[prop-decorator]
@property
def icon_url(self) -> str | None:
return build_icon_url(self.icon_type, self.icon)
@field_validator("updated_at", mode="before")
@classmethod
def _normalize_timestamp(cls, value: datetime | int) -> int:
return to_timestamp(value)
class RecentAppListResponse(ResponseModel):
data: list[RecentAppResponse]
class AppDetail(AppResponseModel):
id: str
name: str
@@ -607,8 +575,6 @@ register_schema_models(
register_response_schema_models(
console_ns,
AppPartial,
RecentAppResponse,
RecentAppListResponse,
AppDetailWithSite,
AppPagination,
)
@@ -733,49 +699,6 @@ class AppListApi(Resource):
return app_detail.model_dump(mode="json"), 201
@console_ns.route("/apps/recent")
class RecentAppListApi(Resource):
@console_ns.doc("list_recent_apps")
@console_ns.doc(description="Get recently modified apps for the home Continue Work section")
@console_ns.doc(params=query_params_from_model(RecentAppListQuery))
@console_ns.response(200, "Success", console_ns.models[RecentAppListResponse.__name__])
@setup_required
@login_required
@account_initialization_required
@enterprise_license_required
@with_session(write=False)
@with_current_user_id
@with_current_tenant_id
def get(self, current_tenant_id: str, current_user_id: str, session: Session):
"""Return the lightweight app cards needed by the Explore home page."""
args = query_params_from_request(RecentAppListQuery)
params = AppListParams(limit=args.limit)
permissions = enterprise_rbac_service.RBACService.MyPermissions.get(
current_tenant_id,
current_user_id,
session=session,
)
if dify_config.RBAC_ENABLED:
access_filter = resolve_app_access_filter(
current_tenant_id,
current_user_id,
session=session,
permissions=permissions,
)
access_filter.apply_to_params(params)
recent_apps = AppService().get_recent_apps(current_user_id, current_tenant_id, params, session)
permission_keys_map = permissions.app.permission_keys_by_resource_ids([app.id for app in recent_apps])
response_items = [
RecentAppResponse.model_validate(app, from_attributes=True).model_copy(
update={"permission_keys": permission_keys_map.get(app.id, [])}
)
for app in recent_apps
]
return dump_response(RecentAppListResponse, {"data": response_items}), 200
@console_ns.route("/apps/starred")
class StarredAppListApi(Resource):
@console_ns.doc("list_starred_apps")
+1 -2
View File
@@ -24,7 +24,7 @@ from extensions.ext_database import db
from fields.base import ResponseModel
from libs.helper import RateLimiter, dump_response, extract_remote_ip, to_timestamp
from models.account import TenantStatus
from models.model import App, AppMode, Site
from models.model import App, Site
from repositories.factory import DifyAPIRepositoryFactory
from services.feature_service import FeatureService
from services.human_input_file_upload_service import HumanInputFileUploadService
@@ -207,7 +207,6 @@ class HumanInputFormApi(Resource):
site=WebAppSiteResponse.from_app_site(
tenant=tenant,
app_model=app_model,
mode=AppMode.value_of(app_model.mode),
site=site,
end_user_id=None,
features=features,
+1 -5
View File
@@ -14,7 +14,7 @@ from extensions.storage.storage_type import StorageType
from fields.base import ResponseModel
from libs.helper import build_icon_url
from models.account import Tenant, TenantStatus
from models.model import App, AppMode, EndUser, IconType, Site
from models.model import App, EndUser, IconType, Site
from services.feature_service import FeatureModel, FeatureService
from services.file_service import FileService
@@ -67,7 +67,6 @@ class WebAppCustomConfigResponse(ResponseModel):
class WebAppSiteResponse(ResponseModel):
app_id: str
mode: AppMode
end_user_id: str | None = None
enable_site: bool
site: WebSiteResponse
@@ -84,7 +83,6 @@ class WebAppSiteResponse(ResponseModel):
*,
tenant: Tenant,
app_model: App,
mode: AppMode,
site: Site,
end_user_id: str | None,
features: FeatureModel,
@@ -111,7 +109,6 @@ class WebAppSiteResponse(ResponseModel):
return cls(
app_id=app_model.id,
mode=mode,
end_user_id=end_user_id,
enable_site=app_model.enable_site,
site=site_response,
@@ -170,7 +167,6 @@ class AppSiteApi(WebApiResource):
return WebAppSiteResponse.from_app_site(
tenant=tenant,
app_model=app_model,
mode=AppMode.value_of(app_model.mode_compatible_with_agent_with_session(session=db.session())),
site=site,
end_user_id=end_user.id,
features=features,
@@ -616,34 +616,23 @@ class AdvancedChatAppGenerator(MessageBasedAppGenerator):
message_snapshot = MessageSnapshot.from_message(message)
session.close()
try:
response = self._handle_advanced_chat_response(
application_generate_entity=application_generate_entity,
workflow=workflow_snapshot,
queue_manager=queue_manager,
conversation=conversation_snapshot,
message=message_snapshot,
user=user,
stream=stream,
draft_var_saver_factory=self._get_draft_var_saver_factory(
invoke_from,
account=user,
tenant_id=application_generate_entity.app_config.tenant_id,
),
)
converted_response = AdvancedChatAppGenerateResponseConverter.convert(
response=response,
invoke_from=invoke_from,
)
except BaseException:
self._join_worker_thread(worker_thread)
raise
# return response or stream generator
response = self._handle_advanced_chat_response(
application_generate_entity=application_generate_entity,
workflow=workflow_snapshot,
queue_manager=queue_manager,
conversation=conversation_snapshot,
message=message_snapshot,
user=user,
stream=stream,
draft_var_saver_factory=self._get_draft_var_saver_factory(
invoke_from,
account=user,
tenant_id=application_generate_entity.app_config.tenant_id,
),
)
if isinstance(converted_response, Generator):
return self._wrap_stream_with_worker_thread_join(converted_response, worker_thread)
self._join_worker_thread(worker_thread)
return converted_response
return AdvancedChatAppGenerateResponseConverter.convert(response=response, invoke_from=invoke_from)
def _generate_worker(
self,
@@ -538,7 +538,6 @@ class AgentAppGenerator(MessageBasedAppGenerator):
request_builder=AgentAppRuntimeRequestBuilder(credentials_provider=credentials_provider),
agent_backend_client=create_agent_backend_run_client(
base_url=dify_config.AGENT_BACKEND_BASE_URL,
api_token=dify_config.AGENT_BACKEND_API_TOKEN,
use_fake=dify_config.AGENT_BACKEND_USE_FAKE,
fake_scenario=dify_config.AGENT_BACKEND_FAKE_SCENARIO,
stream_read_timeout_seconds=dify_config.AGENT_BACKEND_STREAM_READ_TIMEOUT_SECONDS,
-29
View File
@@ -1,5 +1,3 @@
import logging
import threading
from collections.abc import Generator, Mapping, Sequence
from contextlib import AbstractContextManager, nullcontext
from typing import TYPE_CHECKING, Any, Union, final
@@ -25,10 +23,6 @@ from services.workflow_draft_variable_service import DraftVariableSaver as Draft
if TYPE_CHECKING:
from graphon.variables.input_entities import VariableEntity
logger = logging.getLogger(__name__)
_WORKER_THREAD_JOIN_TIMEOUT_SECONDS = 300
@final
class _DebuggerDraftVariableSaver:
@@ -70,29 +64,6 @@ class _DebuggerDraftVariableSaver:
class BaseAppGenerator:
_file_access_controller: DatabaseFileAccessController = DatabaseFileAccessController()
@staticmethod
def _join_worker_thread(worker_thread: threading.Thread) -> None:
# Bound the wait so a leaked app worker cannot occupy an execution slot indefinitely.
worker_thread.join(timeout=_WORKER_THREAD_JOIN_TIMEOUT_SECONDS)
if worker_thread.is_alive():
logger.warning(
"Possible app worker thread leak: thread_name=%s timeout_seconds=%s; "
"continuing without waiting further to avoid occupying an execution slot indefinitely",
worker_thread.name,
_WORKER_THREAD_JOIN_TIMEOUT_SECONDS,
)
@staticmethod
def _wrap_stream_with_worker_thread_join[ResponseT](
response_stream: Generator[ResponseT, None, None],
worker_thread: threading.Thread,
) -> Generator[ResponseT, None, None]:
"""Keep the producer owned by the response stream until both finish."""
try:
yield from response_stream
finally:
BaseAppGenerator._join_worker_thread(worker_thread)
@staticmethod
def _bind_file_access_scope(
*,
@@ -351,28 +351,17 @@ class PipelineGenerator(BaseAppGenerator):
user,
tenant_id=pipeline.tenant_id,
)
try:
response = self._handle_response(
application_generate_entity=application_generate_entity,
workflow=workflow,
queue_manager=queue_manager,
user=user,
stream=streaming,
draft_var_saver_factory=draft_var_saver_factory,
)
converted_response = WorkflowAppGenerateResponseConverter.convert(
response=response,
invoke_from=invoke_from,
)
except BaseException:
self._join_worker_thread(worker_thread)
raise
# return response or stream generator
response = self._handle_response(
application_generate_entity=application_generate_entity,
workflow=workflow,
queue_manager=queue_manager,
user=user,
stream=streaming,
draft_var_saver_factory=draft_var_saver_factory,
)
if isinstance(converted_response, Generator):
return self._wrap_stream_with_worker_thread_join(converted_response, worker_thread)
self._join_worker_thread(worker_thread)
return converted_response
return WorkflowAppGenerateResponseConverter.convert(response=response, invoke_from=invoke_from)
def single_iteration_generate(
self,
+10 -21
View File
@@ -405,28 +405,17 @@ class WorkflowAppGenerator(BaseAppGenerator):
tenant_id=app_model.tenant_id,
)
try:
response = self._handle_response(
application_generate_entity=application_generate_entity,
workflow=workflow,
queue_manager=queue_manager,
user=user,
draft_var_saver_factory=draft_var_saver_factory,
stream=streaming,
)
converted_response = WorkflowAppGenerateResponseConverter.convert(
response=response,
invoke_from=invoke_from,
)
except BaseException:
self._join_worker_thread(worker_thread)
raise
# return response or stream generator
response = self._handle_response(
application_generate_entity=application_generate_entity,
workflow=workflow,
queue_manager=queue_manager,
user=user,
draft_var_saver_factory=draft_var_saver_factory,
stream=streaming,
)
if isinstance(converted_response, Generator):
return self._wrap_stream_with_worker_thread_join(converted_response, worker_thread)
self._join_worker_thread(worker_thread)
return converted_response
return WorkflowAppGenerateResponseConverter.convert(response=response, invoke_from=invoke_from)
def single_iteration_generate(
self,
@@ -39,16 +39,15 @@ class Jinja2TemplateTransformer(TemplateTransformer):
@override
def get_runner_script(cls) -> str:
runner_script = dedent(f"""
import jinja2
import json
from base64 import b64decode
from jinja2.sandbox import SandboxedEnvironment
# declare main function
def main(**inputs):
# Decode base64-encoded template to handle special characters safely
template_code = b64decode('{cls._template_b64_placeholder}').decode('utf-8')
env = SandboxedEnvironment()
template = env.from_string(template_code)
template = jinja2.Template(template_code)
return template.render(**inputs)
# decode and prepare input dict
@@ -68,13 +67,12 @@ class Jinja2TemplateTransformer(TemplateTransformer):
@override
def get_preload_script(cls) -> str:
preload_script = dedent("""
from jinja2.sandbox import SandboxedEnvironment
import jinja2
from base64 import b64decode
def _jinja2_preload_():
# prepare jinja2 sandboxed environment, load template and render
env = SandboxedEnvironment()
template = env.from_string('{{s}}')
# prepare jinja2 environment, load template and render before to avoid sandbox issue
template = jinja2.Template('{{s}}')
template.render(s='a')
if __name__ == '__main__':
+1 -1
View File
@@ -519,7 +519,7 @@ class IndexingRunner:
def filter_string(text):
text = re.sub(r"<\|", "<", text)
text = re.sub(r"\|>", ">", text)
text = re.sub(r"[\x00-\x08\x0B\x0C\x0E-\x1F\x7F]", "", text)
text = re.sub(r"[\x00-\x08\x0B\x0C\x0E-\x1F\x7F\xEF\xBF\xBE]", "", text)
# Unicode U+FFFE
text = re.sub("\ufffe", "", text)
return text
+4 -11
View File
@@ -434,18 +434,14 @@ class PluginService:
exc_info=True,
)
@classmethod
def _fetch_plugin_model_providers_uncached(
cls, tenant_id: str, client: PluginModelClient | None
) -> tuple[ProviderEntity, ...]:
model_client = client or PluginModelClient()
return tuple(cls._to_provider_entity(provider) for provider in model_client.fetch_model_providers(tenant_id))
@classmethod
def _fetch_and_cache_plugin_model_providers(
cls, tenant_id: str, client: PluginModelClient | None, *, refresh_generation: int | None
) -> tuple[ProviderEntity, ...]:
providers = cls._fetch_plugin_model_providers_uncached(tenant_id, client)
model_client = client or PluginModelClient()
providers = tuple(
cls._to_provider_entity(provider) for provider in model_client.fetch_model_providers(tenant_id)
)
generation = cls._load_plugin_model_providers_generation(tenant_id)
if generation is not None and generation == refresh_generation:
cls._store_cached_plugin_model_providers(tenant_id, generation, providers)
@@ -475,9 +471,6 @@ class PluginService:
are intentionally owned by this service so tenant isolation and cache
expiry are handled in one place.
"""
if not dify_config.PLUGIN_MODEL_PROVIDERS_CACHE_ENABLED:
return cls._fetch_plugin_model_providers_uncached(tenant_id, client)
deadline = time.monotonic() + cls.PLUGIN_MODEL_PROVIDERS_LOCK_WAIT_TIMEOUT
while True:
+1 -1
View File
@@ -9,7 +9,7 @@ class CleanProcessor:
# remove invalid symbol
text = re.sub(r"<\|", "<", text)
text = re.sub(r"\|>", ">", text)
text = re.sub(r"[\x00-\x08\x0B\x0C\x0E-\x1F\x7F]", "", text)
text = re.sub(r"[\x00-\x08\x0B\x0C\x0E-\x1F\x7F\xEF\xBF\xBE]", "", text)
# Unicode U+FFFE
text = re.sub("\ufffe", "", text)
+2 -2
View File
@@ -91,8 +91,8 @@ class FixedRecursiveCharacterTextSplitter(EnhanceRecursiveCharacterTextSplitter)
splits = re.split(r" +", text)
else:
splits = text.split(separator)
if self._keep_separator:
splits = [s + separator for s in splits[:-1]] + splits[-1:]
if self._keep_separator:
splits = [s + separator for s in splits[:-1]] + splits[-1:]
else:
splits = list(text)
if separator == "\n":
-1
View File
@@ -497,7 +497,6 @@ class DifyNodeFactory(NodeFactory):
),
"agent_backend_client": create_agent_backend_run_client(
base_url=dify_config.AGENT_BACKEND_BASE_URL,
api_token=dify_config.AGENT_BACKEND_API_TOKEN,
use_fake=dify_config.AGENT_BACKEND_USE_FAKE,
fake_scenario=dify_config.AGENT_BACKEND_FAKE_SCENARIO,
stream_read_timeout_seconds=dify_config.AGENT_BACKEND_STREAM_READ_TIMEOUT_SECONDS,
-27
View File
@@ -5,7 +5,6 @@ from typing import Any
import pytz # type: ignore[import-untyped]
from celery import Celery, Task
from celery.schedules import crontab
from celery.signals import beat_init
from typing_extensions import TypedDict
from configs import dify_config
@@ -37,19 +36,6 @@ class CeleryBeatScheduleEntry(TypedDict):
schedule: crontab | timedelta
def _enqueue_initial_community_telemetry_heartbeat(sender: Any, **_: Any) -> None:
task_name = "community_telemetry.send_heartbeat"
if "community_telemetry_heartbeat" not in sender.app.conf.beat_schedule:
return
task = sender.app.tasks.get(task_name)
if task is not None:
task.apply_async()
beat_init.connect(_enqueue_initial_community_telemetry_heartbeat, weak=False)
def get_celery_ssl_options() -> CelerySSLOptionsDict | None:
"""Get SSL configuration for Celery broker/backend connections."""
# Only apply SSL if we're using Redis as broker/backend
@@ -274,19 +260,6 @@ def init_app(app: DifyApp) -> Celery:
"schedule": timedelta(minutes=dify_config.API_TOKEN_LAST_USED_UPDATE_INTERVAL),
}
if (
dify_config.EDITION == "SELF_HOSTED"
and not dify_config.ENTERPRISE_ENABLED
and not dify_config.DISABLE_TELEMETRY
and not dify_config.DO_NOT_TRACK
and not dify_config.CI
):
imports.append("tasks.community_telemetry_task")
beat_schedule["community_telemetry_heartbeat"] = {
"task": "community_telemetry.send_heartbeat",
"schedule": timedelta(minutes=dify_config.TELEMETRY_HEARTBEAT_INTERVAL_MINUTES),
}
if dify_config.ENTERPRISE_ENABLED and dify_config.ENTERPRISE_TELEMETRY_ENABLED:
imports.append("tasks.enterprise_telemetry_task")
celery_app.conf.update(beat_schedule=beat_schedule, imports=imports)
-1
View File
@@ -383,7 +383,6 @@ class AgentAppComposerResponse(ResponseModel):
variant: Literal[ComposerVariant.AGENT_APP]
agent: AgentComposerAgentResponse
active_config_snapshot: AgentConfigSnapshotSummaryResponse | None = None
active_config_is_published: bool
draft: AgentConfigDraftSummaryResponse | None = None
agent_soul: AgentSoulConfig
save_options: list[ComposerSaveStrategy]
@@ -1,30 +0,0 @@
"""add telemetry fields to dify_setups
Revision ID: 6f5a9c2d8e1b
Revises: d2825e7b9c10
Create Date: 2026-07-23 12:00:00.000000
"""
import sqlalchemy as sa
from alembic import op
# revision identifiers, used by Alembic.
revision = "6f5a9c2d8e1b"
down_revision = "d2825e7b9c10"
branch_labels = None
depends_on = None
def upgrade():
with op.batch_alter_table("dify_setups", schema=None) as batch_op:
batch_op.add_column(sa.Column("instance_id", sa.String(length=255), nullable=True))
batch_op.add_column(sa.Column("install_reported_at", sa.DateTime(), nullable=True))
batch_op.add_column(sa.Column("last_heartbeat_at", sa.DateTime(), nullable=True))
def downgrade():
with op.batch_alter_table("dify_setups", schema=None) as batch_op:
batch_op.drop_column("last_heartbeat_at")
batch_op.drop_column("install_reported_at")
batch_op.drop_column("instance_id")
+2 -5
View File
@@ -362,9 +362,6 @@ class DifySetup(TypeBase):
__table_args__ = (sa.PrimaryKeyConstraint("version", name="dify_setup_pkey"),)
version: Mapped[str] = mapped_column(String(255), nullable=False)
instance_id: Mapped[str | None] = mapped_column(String(255), nullable=True, default=None)
install_reported_at: Mapped[datetime | None] = mapped_column(sa.DateTime, nullable=True, default=None)
last_heartbeat_at: Mapped[datetime | None] = mapped_column(sa.DateTime, nullable=True, default=None)
setup_at: Mapped[datetime] = mapped_column(
sa.DateTime, nullable=False, server_default=func.current_timestamp(), init=False
)
@@ -1117,14 +1114,14 @@ class ExporleBanner(TypeBase):
status: Mapped[BannerStatus] = mapped_column(
EnumText(BannerStatus, length=255),
nullable=False,
server_default=sa.text("'enabled'"),
server_default=sa.text("'enabled'::character varying"),
default=BannerStatus.ENABLED,
)
created_at: Mapped[datetime] = mapped_column(
sa.DateTime, nullable=False, server_default=func.current_timestamp(), init=False
)
language: Mapped[str] = mapped_column(
String(255), nullable=False, server_default=sa.text("'en-US'"), default="en-US"
String(255), nullable=False, server_default=sa.text("'en-US'::character varying"), default="en-US"
)
+1 -40
View File
@@ -1672,23 +1672,6 @@ Create a new application
| 200 | Import confirmed | **application/json**: [Import](#import)<br> |
| 400 | Import failed | **application/json**: [Import](#import)<br> |
### [GET] /apps/recent
**Return the lightweight app cards needed by the Explore home page**
Get recently modified apps for the home Continue Work section
#### Parameters
| Name | Located in | Description | Required | Schema |
| ---- | ---------- | ----------- | -------- | ------ |
| limit | query | Number of recently modified apps to return (1-8) | No | integer, <br>**Default:** 8 |
#### Responses
| Code | Description | Schema |
| ---- | ----------- | ------ |
| 200 | Success | **application/json**: [RecentAppListResponse](#recentapplistresponse)<br> |
### [GET] /apps/starred
Get applications starred by the current account
@@ -13260,7 +13243,6 @@ Model class for AI model.
| Name | Type | Description | Required |
| ---- | ---- | ----------- | -------- |
| active_config_is_published | boolean | | Yes |
| active_config_snapshot | [AgentConfigSnapshotSummaryResponse](#agentconfigsnapshotsummaryresponse) | | No |
| agent | [AgentComposerAgentResponse](#agentcomposeragentresponse) | | Yes |
| agent_soul | [AgentSoulConfig](#agentsoulconfig) | | Yes |
@@ -13300,6 +13282,7 @@ Model class for AI model.
| Name | Type | Description | Required |
| ---- | ---- | ----------- | -------- |
| access_mode | string | | No |
| active_config_is_published | boolean | | No |
| api_base_url | string | | No |
| app_id | string | | No |
| backing_app_id | string | | No |
@@ -21035,28 +21018,6 @@ Whitelist scopes accepted by RBAC app and dataset access config APIs.
| result | string | | Yes |
| updated_at | integer | | Yes |
#### RecentAppListResponse
| Name | Type | Description | Required |
| ---- | ---- | ----------- | -------- |
| data | [ [RecentAppResponse](#recentappresponse) ] | | Yes |
#### RecentAppResponse
| Name | Type | Description | Required |
| ---- | ---- | ----------- | -------- |
| author_name | string | | No |
| icon | string | | No |
| icon_background | string | | No |
| icon_type | [IconType](#icontype) | | No |
| icon_url | string | | Yes |
| id | string | | Yes |
| maintainer | string | | No |
| mode | string, <br>**Available values:** "advanced-chat", "agent-chat", "chat", "completion", "workflow" | *Enum:* `"advanced-chat"`, `"agent-chat"`, `"chat"`, `"completion"`, `"workflow"` | Yes |
| name | string | | Yes |
| permission_keys | [ string ] | | No |
| updated_at | integer | | Yes |
#### RecommendedAppDetailNullableResponse
| Name | Type | Description | Required |
-7
View File
@@ -965,12 +965,6 @@ Returns Server-Sent Events stream.
| ---- | ---- | ----------- | -------- |
| tool_icons | object | Tool icon metadata keyed by tool name | No |
#### AppMode
| Name | Type | Description | Required |
| ---- | ---- | ----------- | -------- |
| AppMode | string | | |
#### AppPermissionQuery
| Name | Type | Description | Required |
@@ -1652,7 +1646,6 @@ in form definition, or a variable while the workflow is running.
| custom_config | [WebAppCustomConfigResponse](#webappcustomconfigresponse) | | No |
| enable_site | boolean | | Yes |
| end_user_id | string | | No |
| mode | [AppMode](#appmode) | | Yes |
| model_config | [WebModelConfigResponse](#webmodelconfigresponse) | | No |
| plan | string | | Yes |
| site | [WebSiteResponse](#websiteresponse) | | Yes |
+1 -7
View File
@@ -75,7 +75,6 @@ from services.errors.account import (
from services.errors.workspace import WorkSpaceNotAllowedCreateError, WorkspacesLimitExceededError
from services.feature_service import FeatureService
from services.plugin.plugin_auto_upgrade_service import PluginAutoUpgradeService
from services.telemetry_service import CommunityTelemetryService
from tasks.delete_account_task import delete_account_task
from tasks.mail_account_deletion_task import send_account_deletion_verification_code
from tasks.mail_change_mail_task import (
@@ -1954,7 +1953,7 @@ class RegisterService:
TenantService.create_owner_tenant_if_not_exist(account=account, is_setup=True, session=session)
dify_setup = DifySetup(version=dify_config.project.version, instance_id=str(uuid.uuid4()))
dify_setup = DifySetup(version=dify_config.project.version)
session.add(dify_setup)
session.commit()
except Exception as e:
@@ -1967,11 +1966,6 @@ class RegisterService:
logger.exception("Setup account failed, email: %s, name: %s", email, name)
raise ValueError(f"Setup failed: {e}")
try:
CommunityTelemetryService.report_install(session=session)
except Exception:
logger.debug("Failed to report install telemetry", exc_info=True)
@classmethod
def register(
cls,
-1
View File
@@ -405,7 +405,6 @@ class AgentComposerService:
"variant": ComposerVariant.AGENT_APP.value,
"agent": cls._serialize_agent(agent),
"active_config_snapshot": cls._serialize_version(version),
"active_config_is_published": bool(agent.active_config_snapshot_id and agent.active_config_is_published),
"draft": cls._serialize_draft(draft),
"agent_soul": draft.config_snapshot_dict,
"save_options": [ComposerSaveStrategy.SAVE_TO_CURRENT_VERSION.value],
-84
View File
@@ -1,7 +1,6 @@
import json
import logging
from collections.abc import Sequence
from dataclasses import dataclass
from datetime import datetime
from typing import Any, Literal, NotRequired, TypedDict, cast, override
@@ -42,20 +41,6 @@ from tasks.remove_app_and_related_data_task import remove_app_and_related_data_t
logger = logging.getLogger(__name__)
AppListSortBy = Literal["last_modified", "recently_created", "earliest_created"]
RecentAppMode = Literal[
AppMode.COMPLETION,
AppMode.WORKFLOW,
AppMode.CHAT,
AppMode.ADVANCED_CHAT,
AppMode.AGENT_CHAT,
]
RECENT_APP_MODES: tuple[RecentAppMode, ...] = (
AppMode.COMPLETION,
AppMode.WORKFLOW,
AppMode.CHAT,
AppMode.ADVANCED_CHAT,
AppMode.AGENT_CHAT,
)
class AppListBaseParams(BaseModel):
@@ -80,19 +65,6 @@ class StarredAppListParams(AppListBaseParams):
pass
@dataclass(frozen=True)
class RecentAppListItem:
id: str
name: str
icon_type: IconType | None
icon: str | None
icon_background: str | None
mode: RecentAppMode
author_name: str | None
updated_at: datetime
maintainer: str | None
class CreateAppParams(BaseModel):
name: str = Field(min_length=1)
description: str | None = None
@@ -351,62 +323,6 @@ class AppService:
return app_models
def get_recent_apps(
self,
user_id: str,
tenant_id: str,
params: AppListParams,
session: Session,
) -> list[RecentAppListItem]:
"""Return recently modified apps as one lightweight, non-paginated projection."""
filters = self._build_app_list_filters(user_id, tenant_id, params, session)
if not filters:
return []
stmt = (
sa.select(
App.id,
App.name,
App.icon_type,
App.icon,
App.icon_background,
App.mode,
Account.name.label("author_name"),
App.updated_at,
App.maintainer,
)
.outerjoin(Account, Account.id == App.created_by)
.where(*filters, App.mode.in_(RECENT_APP_MODES))
.order_by(App.updated_at.desc())
.limit(params.limit)
)
rows = session.execute(stmt).all()
return [
RecentAppListItem(
id=str(app_id),
name=name,
icon_type=icon_type,
icon=icon,
icon_background=icon_background,
mode=cast(RecentAppMode, mode),
author_name=author_name,
updated_at=updated_at,
maintainer=maintainer,
)
for (
app_id,
name,
icon_type,
icon,
icon_background,
mode,
author_name,
updated_at,
maintainer,
) in rows
]
def get_paginate_starred_apps(
self,
user_id: str,
-165
View File
@@ -1,165 +0,0 @@
import logging
import platform
import uuid
from datetime import datetime
from typing import Literal
import httpx
from sqlalchemy import select
from sqlalchemy.orm import Session
from configs import dify_config
from libs.datetime_utils import naive_utc_now
from models.model import DifySetup
logger = logging.getLogger(__name__)
TelemetryEvent = Literal["install", "heartbeat"]
SCHEMA_VERSION = 1
class CommunityTelemetryService:
@classmethod
def report_install(cls, *, session: Session) -> bool:
setup = cls._get_setup(session)
if setup is None:
return False
if setup.instance_id is None:
setup.instance_id = str(uuid.uuid4())
session.add(setup)
session.commit()
payload = cls._build_payload(setup, "install")
if not cls._send_event(payload):
return False
setup.install_reported_at = naive_utc_now()
session.add(setup)
session.commit()
return True
@classmethod
def report_heartbeat(cls, *, session: Session, now: datetime | None = None) -> bool:
setup = cls._get_setup(session)
if setup is None:
return False
if setup.instance_id is None:
setup.instance_id = str(uuid.uuid4())
session.add(setup)
session.commit()
now = now or naive_utc_now()
if not cls._is_heartbeat_due(setup, now):
return False
if setup.install_reported_at is None:
cls.report_install(session=session)
payload = cls._build_payload(setup, "heartbeat")
if not cls._send_event(payload):
return False
setup.last_heartbeat_at = now
session.add(setup)
session.commit()
return True
@classmethod
def _get_setup(cls, session: Session) -> DifySetup | None:
return session.scalar(select(DifySetup).order_by(DifySetup.setup_at.asc()).limit(1))
@classmethod
def _is_enabled(cls) -> bool:
return (
dify_config.EDITION == "SELF_HOSTED"
and not dify_config.ENTERPRISE_ENABLED
and not dify_config.DISABLE_TELEMETRY
and not dify_config.DO_NOT_TRACK
and not dify_config.CI
and bool(dify_config.TELEMETRY_ENDPOINT)
)
@classmethod
def _build_payload(cls, setup: DifySetup, event: TelemetryEvent) -> dict[str, str | int]:
payload: dict[str, str | int] = {
"event": event,
"instance_id": setup.instance_id or "",
"version": setup.version if event == "install" else dify_config.project.version,
"edition": dify_config.EDITION,
"deployment_type": "unknown",
"schema_version": SCHEMA_VERSION,
"os": cls._normalize_os(platform.system()),
"arch": cls._normalize_arch(platform.machine()),
"sent_at": cls._format_datetime(naive_utc_now()),
}
if event == "install":
payload["installed_at"] = cls._format_datetime(setup.setup_at)
return payload
@classmethod
def _send_event(cls, payload: dict[str, str | int]) -> bool:
if not cls._is_enabled():
return False
endpoints = [dify_config.TELEMETRY_ENDPOINT]
if dify_config.TELEMETRY_FALLBACK_ENDPOINT not in endpoints:
endpoints.append(dify_config.TELEMETRY_FALLBACK_ENDPOINT)
for endpoint in endpoints:
if not endpoint:
continue
try:
response = httpx.post(
endpoint,
json=payload,
timeout=dify_config.TELEMETRY_TIMEOUT_SECONDS,
)
response.raise_for_status()
return True
except httpx.RequestError:
logger.debug("Failed to send community telemetry event to %s", endpoint, exc_info=True)
except httpx.HTTPStatusError:
logger.debug("Community telemetry endpoint returned an error: %s", endpoint, exc_info=True)
return False
return False
@classmethod
def _is_heartbeat_due(cls, setup: DifySetup, now: datetime) -> bool:
if setup.instance_id is None:
return False
if setup.last_heartbeat_at is not None and setup.last_heartbeat_at.date() >= now.date():
return False
return True
@staticmethod
def _format_datetime(value: datetime) -> str:
return value.replace(microsecond=0).isoformat() + "Z"
@staticmethod
def _normalize_os(value: str) -> str:
os_name = value.lower()
if os_name in {"linux", "darwin", "windows"}:
return os_name
return "unknown"
@staticmethod
def _normalize_arch(value: str) -> str:
arch = value.lower()
if arch in {"x86_64", "amd64"}:
return "amd64"
if arch in {"aarch64", "arm64"}:
return "arm64"
if arch.startswith("arm"):
return "arm"
if arch in {"i386", "i686", "x86"}:
return "386"
return "unknown"
+9 -9
View File
@@ -278,14 +278,14 @@ class VariableTruncator(BaseTruncator):
target_length = self._array_element_limit
for i, item in enumerate(value):
# ``File`` is routed through ``_truncate_json_primitives`` (whose
# dedicated ``File`` branch returns the file as-is with its real
# serialized size). That preserves the count cap
# (``array_element_limit``) and the byte budget (``target_size``)
# for ``list[File]`` — the original "Dirty fix" branch above this
# loop bypassed both guarantees and reported ``used_size=2`` even
# when the returned array serialized to well over the budget.
# See https://github.com/langgenius/dify/issues/39218.
# Dirty fix:
# The output of `Start` node may contain list of `File` elements,
# causing `AssertionError` while invoking `_truncate_json_primitives`.
#
# This check ensures that `list[File]` are handled separately
if isinstance(item, File):
truncated_value.append(item)
continue
if i >= target_length:
return _PartResult(truncated_value, used_size, True)
if i > 0:
@@ -295,7 +295,7 @@ class VariableTruncator(BaseTruncator):
break
remaining_budget = target_size - used_size
if item is None or isinstance(item, (str, list, dict, bool, int, float, File, UpdatedVariable)):
if item is None or isinstance(item, (str, list, dict, bool, int, float, UpdatedVariable)):
part_result = self._truncate_json_primitives(item, remaining_budget)
else:
raise UnknownTypeError(f"got unknown type {type(item)} in array truncation")
+1 -1
View File
@@ -116,7 +116,7 @@ class WebAppAuthService:
@classmethod
def _get_account_jwt_token(cls, account: Account) -> str:
exp_dt = datetime.now(UTC) + timedelta(minutes=dify_config.ACCESS_TOKEN_EXPIRE_MINUTES)
exp_dt = datetime.now(UTC) + timedelta(minutes=dify_config.ACCESS_TOKEN_EXPIRE_MINUTES * 24)
exp = int(exp_dt.timestamp())
payload = {
+1 -1
View File
@@ -27,6 +27,7 @@ class WorkspaceService:
tenant_info: dict[str, object] = {
"id": tenant.id,
"name": tenant.name,
"plan": tenant.plan,
"status": tenant.status,
"created_at": tenant.created_at,
"trial_end_reason": None,
@@ -43,7 +44,6 @@ class WorkspaceService:
tenant_info["role"] = tenant_account_join.role
feature = FeatureService.get_features(tenant.id, exclude_vector_space=True)
tenant_info["plan"] = feature.billing.subscription.plan if feature.billing.enabled else None
can_replace_logo = feature.can_replace_logo
if can_replace_logo and TenantService.has_roles(
@@ -22,7 +22,6 @@ def _create_agent_backend_client():
return None
return create_agent_backend_run_client(
base_url=dify_config.AGENT_BACKEND_BASE_URL,
api_token=dify_config.AGENT_BACKEND_API_TOKEN,
use_fake=dify_config.AGENT_BACKEND_USE_FAKE,
fake_scenario=dify_config.AGENT_BACKEND_FAKE_SCENARIO,
stream_read_timeout_seconds=dify_config.AGENT_BACKEND_STREAM_READ_TIMEOUT_SECONDS,
@@ -457,7 +457,7 @@ def _publish_streaming_response(
@shared_task(queue=WORKFLOW_BASED_APP_EXECUTION_QUEUE)
def workflow_based_app_execution_task(
payload: str,
) -> Mapping[str, Any] | None:
) -> Generator[Mapping[str, Any] | str, None, None] | Mapping[str, Any] | None:
exec_params = AppExecutionParams.model_validate_json(payload)
logger.info("workflow_based_app_execution_task run with params: %s", exec_params)
-19
View File
@@ -1,19 +0,0 @@
import logging
from celery import shared_task
from sqlalchemy.orm import sessionmaker
from extensions.ext_database import db
from services.telemetry_service import CommunityTelemetryService
logger = logging.getLogger(__name__)
@shared_task(name="community_telemetry.send_heartbeat", queue="schedule_executor")
def send_community_telemetry_heartbeat() -> None:
session_factory = sessionmaker(bind=db.engine, expire_on_commit=False)
with session_factory() as session:
try:
CommunityTelemetryService.report_heartbeat(session=session)
except Exception:
logger.debug("Failed to process community telemetry heartbeat", exc_info=True)
@@ -17,11 +17,10 @@ class _SimpleJinja2Renderer:
"""Minimal Jinja2-based renderer for integration tests (no code executor)."""
def render_template(self, template: str, variables: dict[str, object]) -> str:
from jinja2.sandbox import SandboxedEnvironment
from jinja2 import Template
try:
env = SandboxedEnvironment()
return env.from_string(template).render(**variables)
return Template(template).render(**variables)
except Exception as exc:
raise TemplateRenderError(str(exc)) from exc
@@ -97,7 +97,6 @@ class TestAppSiteApi:
assert result["end_user_id"] == end_user.id
assert result["plan"] == "basic"
assert result["enable_site"] is True
assert result["mode"] == AppMode.CHAT
@patch("controllers.web.site.FileService.get_file_presigned_url")
@patch("controllers.web.site.FeatureService.get_features")
@@ -179,7 +178,6 @@ class TestWebAppSiteResponse:
response = WebAppSiteResponse.from_app_site(
tenant=tenant,
app_model=app_model,
mode=AppMode.CHAT,
site=_site_model(app_id=app_model.id),
end_user_id="eu-1",
features=FeatureModel(can_replace_logo=False, webapp_copyright_enabled=True),
@@ -187,7 +185,6 @@ class TestWebAppSiteResponse:
)
assert response.app_id == app_model.id
assert response.mode == AppMode.CHAT
assert response.end_user_id == "eu-1"
assert response.enable_site is True
assert response.plan == "basic"
@@ -212,7 +209,6 @@ class TestWebAppSiteResponse:
response = WebAppSiteResponse.from_app_site(
tenant=tenant,
app_model=app_model,
mode=AppMode.CHAT,
site=site,
end_user_id=None,
features=FeatureModel(can_replace_logo=False, webapp_copyright_enabled=True),
@@ -240,7 +236,6 @@ class TestWebAppSiteResponse:
response = WebAppSiteResponse.from_app_site(
tenant=tenant,
app_model=app_model,
mode=AppMode.CHAT,
site=_site_model(app_id=app_model.id),
end_user_id="eu-1",
features=FeatureModel(can_replace_logo=True, webapp_copyright_enabled=True),
@@ -24,10 +24,7 @@ class TestWorkspaceService:
patch("services.workspace_service.dify_config") as mock_dify_config,
):
# Setup default mock returns
feature = mock_feature_service.get_features.return_value
feature.can_replace_logo = True
feature.billing.enabled = True
feature.billing.subscription.plan = "professional"
mock_feature_service.get_features.return_value.can_replace_logo = True
mock_tenant_service.has_roles.return_value = True
mock_dify_config.FILES_URL = "https://example.com/files"
@@ -115,7 +112,7 @@ class TestWorkspaceService:
assert result is not None
assert result["id"] == tenant.id
assert result["name"] == tenant.name
assert result["plan"] == "professional"
assert result["plan"] == tenant.plan
assert result["status"] == tenant.status
assert result["role"] == TenantAccountRole.OWNER
assert result["created_at"] == tenant.created_at
@@ -162,7 +159,7 @@ class TestWorkspaceService:
assert result is not None
assert result["id"] == tenant.id
assert result["name"] == tenant.name
assert result["plan"] == "professional"
assert result["plan"] == tenant.plan
assert result["status"] == tenant.status
assert result["role"] == TenantAccountRole.OWNER
assert result["created_at"] == tenant.created_at
@@ -217,7 +214,7 @@ class TestWorkspaceService:
assert result is not None
assert result["id"] == tenant.id
assert result["name"] == tenant.name
assert result["plan"] == "professional"
assert result["plan"] == tenant.plan
assert result["status"] == tenant.status
assert result["role"] == TenantAccountRole.NORMAL
assert result["created_at"] == tenant.created_at
@@ -609,23 +606,20 @@ class TestWorkspaceService:
def test_get_tenant_info_should_not_include_cloud_fields_in_self_hosted(
self, db_session_with_containers: Session, mock_external_service_dependencies
):
"""Cloud-only billing data should not appear in SELF_HOSTED mode."""
"""next_credit_reset_date and trial_credits should NOT appear in SELF_HOSTED mode."""
fake = Faker()
account, tenant = self._create_test_account_and_tenant(
db_session_with_containers, mock_external_service_dependencies
)
mock_external_service_dependencies["dify_config"].DEPLOYMENT_EDITION = DeploymentEdition.COMMUNITY
feature = mock_external_service_dependencies["feature_service"].get_features.return_value
feature.can_replace_logo = False
feature.billing.enabled = False
mock_external_service_dependencies["feature_service"].get_features.return_value.can_replace_logo = False
mock_external_service_dependencies["tenant_service"].has_roles.return_value = False
with patch("services.workspace_service.current_user", account):
result = WorkspaceService.get_tenant_info(tenant, db_session_with_containers)
assert result is not None
assert result["plan"] is None
assert "next_credit_reset_date" not in result
assert "trial_credits" not in result
assert "trial_credits_used" not in result
@@ -4,7 +4,6 @@ from dotenv import dotenv_values
BASE_API_AND_DOCKER_CONFIG_SET_DIFF: frozenset[str] = frozenset(
(
"AGENT_BACKEND_API_TOKEN",
"APP_MAX_EXECUTION_TIME",
"BATCH_UPLOAD_LIMIT",
"CELERY_BEAT_SCHEDULER_TIME",
@@ -44,7 +43,6 @@ BASE_API_AND_DOCKER_CONFIG_SET_DIFF: frozenset[str] = frozenset(
BASE_API_AND_DOCKER_COMPOSE_CONFIG_SET_DIFF: frozenset[str] = frozenset(
(
"AGENT_BACKEND_API_TOKEN",
"BATCH_UPLOAD_LIMIT",
"CELERY_BEAT_SCHEDULER_TIME",
"HTTP_REQUEST_MAX_CONNECT_TIMEOUT",
+26 -56
View File
@@ -1,13 +1,11 @@
import os
import shutil
from collections.abc import Iterator
from pathlib import Path
from unittest.mock import MagicMock, patch
import pytest
from flask import Flask
from sqlalchemy import create_engine
from sqlalchemy.engine import URL, Engine
from sqlalchemy.engine import Engine
from sqlalchemy.orm import Session, sessionmaker
# Getting the absolute path of the current file's directory
@@ -37,7 +35,7 @@ os.environ.setdefault("OPENDAL_SCHEME", "fs")
os.environ.setdefault("OPENDAL_FS_ROOT", "/tmp/dify-storage")
os.environ.setdefault("STORAGE_TYPE", "opendal")
import core.db.session_factory as session_factory_module
from core.db.session_factory import configure_session_factory, session_factory
from extensions import ext_redis
from models.account import Account, Tenant, TenantAccountJoin, TenantAccountRole
from models.base import TypeBase
@@ -113,70 +111,42 @@ def reset_secret_key():
dify_config.SECRET_KEY = original
@pytest.fixture(scope="session")
def _unit_test_engine():
engine = create_engine("sqlite:///:memory:")
yield engine
engine.dispose()
@pytest.fixture
def _sqlite_engine(_sqlite_database_template: Path, tmp_path: Path) -> Iterator[Engine]:
"""Create an engine over a pristine per-test copy of the SQLite schema."""
database_path = tmp_path / "unit-tests.sqlite3"
shutil.copyfile(_sqlite_database_template, database_path)
engine = create_engine(URL.create("sqlite", database=str(database_path)))
def sqlite_engine() -> Iterator[Engine]:
"""Create an isolated in-memory SQLite engine for tests that need a disposable database."""
engine = create_engine("sqlite:///:memory:")
try:
yield engine
finally:
engine.dispose()
database_path.unlink(missing_ok=True)
@pytest.fixture(scope="session")
def _sqlite_database_template(tmp_path_factory: pytest.TempPathFactory) -> Path:
"""Create one empty full-schema SQLite database per pytest worker."""
@pytest.fixture
def sqlite_session(request: pytest.FixtureRequest, sqlite_engine: Engine) -> Iterator[Session]:
"""Yield a SQLite session after creating the model tables passed through ``request.param``."""
database_path = tmp_path_factory.mktemp("sqlite-template") / "unit-tests.sqlite3"
engine = create_engine(URL.create("sqlite", database=str(database_path)))
try:
TypeBase.metadata.create_all(engine)
finally:
engine.dispose()
return database_path
models: tuple[type[TypeBase], ...] = request.param
tables = [model.metadata.tables[model.__tablename__] for model in models]
TypeBase.metadata.create_all(sqlite_engine, tables=tables)
session_factory = sessionmaker(bind=sqlite_engine, expire_on_commit=False)
with session_factory() as session:
yield session
@pytest.fixture(autouse=True)
def _sqlite_session_factory(
_sqlite_engine: Engine,
monkeypatch: pytest.MonkeyPatch,
) -> sessionmaker[Session]:
"""Bind all unit-test Sessions to the pristine full-schema SQLite database."""
factory = sessionmaker(bind=_sqlite_engine, expire_on_commit=False)
monkeypatch.setattr(session_factory_module, "_session_maker", factory)
return factory
@pytest.fixture
def sqlite_engine(_sqlite_engine: Engine) -> Engine:
"""Expose the pristine full-schema SQLite engine to tests."""
return _sqlite_engine
@pytest.fixture
def sqlite_session_factory(_sqlite_session_factory: sessionmaker[Session]) -> sessionmaker[Session]:
"""Expose the shared SQLite session factory to tests."""
return _sqlite_session_factory
@pytest.fixture
def sqlite_session(_sqlite_session_factory: sessionmaker[Session]) -> Iterator[Session]:
"""Yield a session over the pristine full-schema SQLite database.
Legacy indirect model parameters remain accepted by pytest but are ignored.
Remove those decorators as their test files receive individual review.
"""
with _sqlite_session_factory() as session:
yield session
def _configure_session_factory(_unit_test_engine):
try:
session_factory.get_session_maker()
except RuntimeError:
configure_session_factory(_unit_test_engine, expire_on_commit=False)
def persist_service_api_tenant_owner(session: Session, tenant: Tenant, owner: Account) -> TenantAccountJoin:
@@ -115,7 +115,6 @@ def _agent_app_composer_response() -> dict:
"active_config_snapshot_id": "version-1",
},
"active_config_snapshot": _version_response(),
"active_config_is_published": True,
"agent_soul": {},
"save_options": ["save_to_current_version"],
}
@@ -377,7 +376,7 @@ def test_agent_app_list_and_create_use_agent_route(
assert created["app_id"] == "app-created"
assert created["debug_conversation_id"] == "debug-conversation-created"
assert created["role"] == "Created role"
assert "active_config_is_published" not in created
assert created["active_config_is_published"] is False
assert "bound_agent_id" not in created
create_call = cast(dict[str, object], captured["create"])
create_params = cast(Any, create_call["params"])
@@ -488,7 +487,7 @@ def test_agent_app_detail_update_delete_resolve_app_from_agent_id(
assert detail["debug_conversation_has_messages"] is True
assert detail["debug_conversation_message_count"] == 2
assert detail["role"] == "Resolved role"
assert "active_config_is_published" not in detail
assert detail["active_config_is_published"] is False
assert "bound_agent_id" not in detail
assert captured["get_app"] == {"app": app_model, "session": session}
with app.test_request_context(
@@ -503,7 +502,7 @@ def test_agent_app_detail_update_delete_resolve_app_from_agent_id(
assert updated["debug_conversation_has_messages"] is True
assert updated["debug_conversation_message_count"] == 2
assert updated["role"] == "Resolved role"
assert "active_config_is_published" not in updated
assert updated["active_config_is_published"] is False
assert "bound_agent_id" not in updated
update_call = cast(dict[str, object], captured["update"])
assert update_call["app"] is app_model
@@ -846,6 +845,9 @@ def test_agent_app_update_allows_empty_role(app: Flask, monkeypatch: pytest.Monk
monkeypatch.setattr(
roster_controller.AgentRosterService, "count_agent_app_debug_conversation_messages", lambda _self, **kwargs: 0
)
monkeypatch.setattr(
roster_controller.AgentRosterService, "active_config_is_published", lambda _self, **kwargs: False
)
monkeypatch.setattr(
roster_controller.FeatureService,
"get_system_features",
@@ -1297,14 +1299,13 @@ def test_agent_composer_routes_resolve_app_from_agent_id(
composer_controller.AgentComposerService, "collect_validation_findings", collect_validation_findings
)
monkeypatch.setattr(composer_controller.AgentComposerService, "get_agent_app_candidates", get_agent_app_candidates)
composer = unwrap(AgentComposerApi.get)(AgentComposerApi(), MagicMock(), "tenant-1", agent_id)
assert composer["variant"] == "agent_app"
assert composer["active_config_is_published"] is True
assert unwrap(AgentComposerApi.get)(AgentComposerApi(), MagicMock(), "tenant-1", agent_id)["variant"] == "agent_app"
assert cast(dict[str, object], captured["load"])["agent_id"] == agent_id
with app.test_request_context(json=payload):
saved_composer = unwrap(AgentComposerApi.put)(AgentComposerApi(), MagicMock(), "tenant-1", account_id, agent_id)
assert saved_composer["variant"] == "agent_app"
assert saved_composer["active_config_is_published"] is True
assert (
unwrap(AgentComposerApi.put)(AgentComposerApi(), MagicMock(), "tenant-1", account_id, agent_id)["variant"]
== "agent_app"
)
assert cast(dict[str, object], captured["save"])["agent_id"] == agent_id
assert unwrap(AgentComposerValidateApi.post)(AgentComposerValidateApi(), MagicMock(), "tenant-1", agent_id) == {
"result": "success",
@@ -708,121 +708,6 @@ def test_app_list_api_attaches_permission_keys(app, app_module):
assert resp["data"][0]["permission_keys"] == ["app.acl.view_layout", "app.acl.edit"]
def test_recent_app_list_api_returns_only_home_card_fields(app, app_module):
method = app_module.RecentAppListApi.get
while hasattr(method, "__wrapped__"):
method = method.__wrapped__
recent_app = SimpleNamespace(
id="app-1",
name="Recent App",
icon_type="emoji",
icon="🚀",
icon_background="#FFFFFF",
mode="chat",
author_name="Recent Author",
updated_at=_ts(15),
maintainer="acct-1",
)
get_recent_apps = MagicMock(return_value=[recent_app])
with app.test_request_context("/apps/recent?limit=8"):
with pytest.MonkeyPatch.context() as monkeypatch:
monkeypatch.setattr(dify_config, "RBAC_ENABLED", False)
monkeypatch.setattr(app_module.AppService, "get_recent_apps", get_recent_apps)
monkeypatch.setattr(
app_module.enterprise_rbac_service.RBACService.MyPermissions,
"get",
lambda tenant_id, account_id, session: app_module.enterprise_rbac_service.MyPermissionsResponse(
app=app_module.enterprise_rbac_service.ResourcePermissionSnapshot(
overrides=[
app_module.enterprise_rbac_service.ResourcePermissionKeys(
resource_id="app-1",
permission_keys=["app.acl.monitor"],
)
]
)
),
)
resp, status = method(app_module.RecentAppListApi(), "tenant-1", "acct-1", MagicMock())
assert status == 200
assert resp == {
"data": [
{
"id": "app-1",
"name": "Recent App",
"icon_type": "emoji",
"icon": "🚀",
"icon_background": "#FFFFFF",
"mode": "chat",
"author_name": "Recent Author",
"updated_at": int(_ts(15).timestamp()),
"permission_keys": ["app.acl.monitor"],
"maintainer": "acct-1",
"icon_url": None,
}
]
}
params = get_recent_apps.call_args.args[2]
assert params.limit == 8
assert "total" not in resp
assert "description" not in resp["data"][0]
assert "tags" not in resp["data"][0]
assert "workflow" not in resp["data"][0]
@pytest.mark.parametrize("mode", ["channel", "rag-pipeline", "agent"])
def test_recent_app_response_rejects_non_home_app_modes(app_module, mode: str) -> None:
with pytest.raises(ValidationError):
app_module.RecentAppResponse.model_validate(
{
"id": "app-1",
"name": "Recent App",
"mode": mode,
"updated_at": _ts(),
}
)
def test_recent_app_list_api_applies_rbac_visibility_filter(app, app_module):
method = app_module.RecentAppListApi.get
while hasattr(method, "__wrapped__"):
method = method.__wrapped__
get_recent_apps = MagicMock(return_value=[])
with app.test_request_context("/apps/recent"):
with pytest.MonkeyPatch.context() as monkeypatch:
monkeypatch.setattr(dify_config, "RBAC_ENABLED", True)
monkeypatch.setattr(app_module.AppService, "get_recent_apps", get_recent_apps)
monkeypatch.setattr(
app_module.enterprise_rbac_service.RBACService.MyPermissions,
"get",
lambda tenant_id, account_id, session: app_module.enterprise_rbac_service.MyPermissionsResponse(
workspace=app_module.enterprise_rbac_service.WorkspacePermissionSnapshot(
permission_keys=["app.create_and_management"]
)
),
)
monkeypatch.setattr(
app_module.enterprise_rbac_service.RBACService.AppAccess,
"whitelist_resources",
lambda tenant_id, account_id: SimpleNamespace(
unrestricted=False,
resource_ids=["app-shared"],
),
)
resp, status = method(app_module.RecentAppListApi(), "tenant-1", "acct-1", MagicMock())
assert status == 200
assert resp == {"data": []}
params = get_recent_apps.call_args.args[2]
assert params.accessible_app_ids == ["app-shared"]
assert params.include_own_apps is True
def test_app_list_api_limits_to_apps_created_by_current_user_without_view_permission(app, app_module):
method = app_module.AppListApi.get
while hasattr(method, "__wrapped__"):
@@ -163,7 +163,6 @@ def test_get_form_includes_site(monkeypatch: pytest.MonkeyPatch, app: Flask, dat
assert body["expiration_time"] == int(expiration_time.timestamp())
assert body["site"] == {
"app_id": app_model.id,
"mode": "chat",
"end_user_id": None,
"enable_site": True,
"site": {
@@ -384,7 +383,6 @@ def test_get_form_allows_backstage_token(monkeypatch: pytest.MonkeyPatch, app: F
assert body["expiration_time"] == int(expiration_time.timestamp())
assert body["site"] == {
"app_id": app_model.id,
"mode": "chat",
"end_user_id": None,
"enable_site": True,
"site": {
@@ -3,42 +3,7 @@ from unittest.mock import MagicMock, patch
from configs import dify_config
from controllers.web import site as site_module
from extensions.storage.storage_type import StorageType
from models.model import AppMode, IconType, Site
from services.feature_service import FeatureModel
def test_app_site_api_returns_legacy_agent_compatible_mode() -> None:
app_model = MagicMock()
app_model.id = "app-id"
app_model.tenant_id = "tenant-id"
app_model.tenant = MagicMock(id="tenant-id", status="normal")
app_model.mode_compatible_with_agent_with_session.return_value = AppMode.AGENT_CHAT
end_user = MagicMock(id="end-user-id")
site = MagicMock(spec=Site)
response = MagicMock()
response.model_dump.return_value = {"mode": AppMode.AGENT_CHAT}
with (
patch.object(site_module, "db") as mock_db,
patch.object(site_module.FeatureService, "get_features", return_value=FeatureModel(can_replace_logo=False)),
patch.object(site_module, "_build_site_icon_url", return_value=None),
patch.object(site_module.WebAppSiteResponse, "from_app_site", return_value=response) as mock_from_app_site,
):
mock_db.session.scalar.return_value = site
result = site_module.AppSiteApi().get(app_model, end_user)
assert result["mode"] == AppMode.AGENT_CHAT
app_model.mode_compatible_with_agent_with_session.assert_called_once_with(session=mock_db.session())
mock_from_app_site.assert_called_once_with(
tenant=app_model.tenant,
app_model=app_model,
mode=AppMode.AGENT_CHAT,
site=site,
end_user_id=end_user.id,
features=FeatureModel(can_replace_logo=False),
can_replace_logo=False,
icon_url=None,
)
from models.model import IconType, Site
def test_build_site_icon_url_uses_s3_presigned_url() -> None:
@@ -442,13 +442,6 @@ class TestAdvancedChatAppGeneratorInternals:
def start(self):
thread_data["started"] = True
def join(self, timeout):
thread_data["joined"] = True
thread_data["join_timeout"] = timeout
def is_alive(self):
return False
monkeypatch.setattr("core.app.apps.advanced_chat.app_generator.threading.Thread", _Thread)
monkeypatch.setattr(
"core.app.apps.advanced_chat.app_generator.db", SimpleNamespace(engine=object(), session=db_session)
@@ -482,8 +475,6 @@ class TestAdvancedChatAppGeneratorInternals:
assert response["response"] == {"raw": True}
assert thread_data["started"] is True
assert thread_data["joined"] is True
assert thread_data["join_timeout"] == 300
assert "pause-layer" in thread_data["kwargs"]["graph_engine_layers"]
assert generator._dialogue_count == 3
assert init_records.call_args.kwargs["session"] is db_session
@@ -551,13 +542,6 @@ class TestAdvancedChatAppGeneratorInternals:
def start(self):
thread_data["started"] = True
def join(self, timeout):
thread_data["joined"] = True
thread_data["join_timeout"] = timeout
def is_alive(self):
return False
monkeypatch.setattr("core.app.apps.advanced_chat.app_generator.threading.Thread", _Thread)
monkeypatch.setattr(
"core.app.apps.advanced_chat.app_generator.db", SimpleNamespace(engine=object(), session=db_session)
@@ -590,8 +574,6 @@ class TestAdvancedChatAppGeneratorInternals:
init_records.assert_not_called()
get_thread_messages_length.assert_called_once_with(conversation.id, session=db_session)
assert thread_data["started"] is True
assert thread_data["joined"] is True
assert thread_data["join_timeout"] == 300
db_session.commit.assert_not_called()
db_session.refresh.assert_not_called()
db_session.close.assert_called_once()
@@ -435,7 +435,6 @@ def test_generate_success_returns_converted(generator, mocker: MockerFixture):
mocker.patch.object(module, "PipelineQueueManager", return_value=queue_manager)
worker_thread = MagicMock()
worker_thread.is_alive.return_value = False
mocker.patch.object(module.threading, "Thread", return_value=worker_thread)
mocker.patch.object(generator, "_get_draft_var_saver_factory", return_value=MagicMock())
@@ -462,7 +461,6 @@ def test_generate_success_returns_converted(generator, mocker: MockerFixture):
)
assert result == "converted"
worker_thread.join.assert_called_once_with(timeout=300)
def test_single_iteration_generate_validates_inputs(generator, mocker: MockerFixture):
@@ -1,6 +1,3 @@
import logging
from unittest.mock import Mock
import pytest
from core.app.apps.base_app_generator import BaseAppGenerator
@@ -372,58 +369,6 @@ def test_validate_inputs_optional_file_with_empty_string_ignores_default():
class TestBaseAppGeneratorExtras:
def test_wrap_stream_joins_worker_after_stream_exhaustion(self):
base_app_generator = BaseAppGenerator()
worker_thread = Mock()
worker_thread.is_alive.return_value = False
def response_stream():
yield {"event": "workflow_finished"}
managed_stream = base_app_generator._wrap_stream_with_worker_thread_join(
response_stream(),
worker_thread,
)
assert next(managed_stream) == {"event": "workflow_finished"}
worker_thread.join.assert_not_called()
with pytest.raises(StopIteration):
next(managed_stream)
worker_thread.join.assert_called_once_with(timeout=300)
def test_wrap_stream_joins_worker_when_stream_closes(self):
base_app_generator = BaseAppGenerator()
worker_thread = Mock()
worker_thread.is_alive.return_value = False
def response_stream():
yield {"event": "workflow_started"}
yield {"event": "workflow_finished"}
managed_stream = base_app_generator._wrap_stream_with_worker_thread_join(
response_stream(),
worker_thread,
)
assert next(managed_stream) == {"event": "workflow_started"}
managed_stream.close()
worker_thread.join.assert_called_once_with(timeout=300)
def test_join_worker_thread_warns_when_thread_remains_alive(self, caplog: pytest.LogCaptureFixture):
worker_thread = Mock()
worker_thread.name = "leaked-app-worker"
worker_thread.is_alive.return_value = True
with caplog.at_level(logging.WARNING, logger="core.app.apps.base_app_generator"):
BaseAppGenerator._join_worker_thread(worker_thread)
worker_thread.join.assert_called_once_with(timeout=300)
assert "Possible app worker thread leak" in caplog.text
assert "leaked-app-worker" in caplog.text
def test_prepare_user_inputs_converts_files_and_lists(self, monkeypatch: pytest.MonkeyPatch):
base_app_generator = BaseAppGenerator()
@@ -211,13 +211,6 @@ def test_generate_appends_pause_layer_and_forwards_state(mocker: MockerFixture):
def start(self):
return None
def join(self, timeout):
worker_kwargs["joined"] = True
worker_kwargs["join_timeout"] = timeout
def is_alive(self):
return False
mocker.patch("core.app.apps.workflow.app_generator.threading.Thread", DummyThread)
app_model = SimpleNamespace(mode="workflow", tenant_id="tenant")
@@ -251,8 +244,6 @@ def test_generate_appends_pause_layer_and_forwards_state(mocker: MockerFixture):
assert result == "converted"
assert worker_kwargs["kwargs"]["graph_engine_layers"] == ("base-layer", pause_layer)
assert worker_kwargs["kwargs"]["graph_runtime_state"] is graph_runtime_state
assert worker_kwargs["joined"] is True
assert worker_kwargs["join_timeout"] == 300
assert draft_saver_factory.call_args.kwargs["tenant_id"] == app_model.tenant_id
@@ -295,8 +286,6 @@ def test_resume_path_runs_worker_with_runtime_state(mocker: MockerFixture):
mocker.patch("core.app.apps.workflow.app_generator.WorkflowAppRunner", side_effect=runner_ctor)
worker_lifecycle: dict[str, bool] = {}
class ImmediateThread:
def __init__(self, target, kwargs):
target(**kwargs)
@@ -304,13 +293,6 @@ def test_resume_path_runs_worker_with_runtime_state(mocker: MockerFixture):
def start(self):
return None
def join(self, timeout):
worker_lifecycle["joined"] = True
worker_lifecycle["join_timeout"] = timeout
def is_alive(self):
return False
mocker.patch("core.app.apps.workflow.app_generator.threading.Thread", ImmediateThread)
mocker.patch(
@@ -349,7 +331,5 @@ def test_resume_path_runs_worker_with_runtime_state(mocker: MockerFixture):
)
assert result == "raw-response"
assert worker_lifecycle["joined"] is True
assert worker_lifecycle["join_timeout"] == 300
runner_instance.run.assert_called_once()
queue_manager.graph_runtime_state = runtime_state
@@ -1,9 +1,5 @@
import threading
from collections.abc import Generator
import pytest
from core.app.apps.base_app_generator import BaseAppGenerator
from core.app.apps.workflow.active_workflow_tasks import (
active_workflow_task,
get_active_workflow_task_count,
@@ -32,51 +28,3 @@ def test_active_workflow_task_rejects_duplicate_task_id() -> None:
with pytest.raises(ValueError, match="already active"):
with active_workflow_task("task-a"):
pass
def test_managed_stream_waits_for_active_worker_cleanup() -> None:
worker_started = threading.Event()
release_worker = threading.Event()
stream_exhausted = threading.Event()
consumer_finished = threading.Event()
consumer_errors: list[BaseException] = []
def run_worker() -> None:
with active_workflow_task("task-a"):
worker_started.set()
release_worker.wait()
def response_stream() -> Generator[dict[str, str], None, None]:
yield {"event": "workflow_finished"}
stream_exhausted.set()
worker_thread = threading.Thread(target=run_worker)
worker_thread.start()
assert worker_started.wait(timeout=2)
managed_stream = BaseAppGenerator._wrap_stream_with_worker_thread_join(response_stream(), worker_thread)
assert next(managed_stream) == {"event": "workflow_finished"}
def finish_stream() -> None:
try:
list(managed_stream)
except BaseException as exc:
consumer_errors.append(exc)
finally:
consumer_finished.set()
consumer_thread = threading.Thread(target=finish_stream)
consumer_thread.start()
try:
assert stream_exhausted.wait(timeout=2)
assert not consumer_finished.is_set()
assert get_active_workflow_task_count() == 1
finally:
release_worker.set()
consumer_thread.join(timeout=2)
worker_thread.join(timeout=2)
assert not consumer_thread.is_alive()
assert not worker_thread.is_alive()
assert consumer_errors == []
assert get_active_workflow_task_count() == 0
@@ -15,70 +15,6 @@ from models.model import AppMode
class TestWorkflowAppGeneratorValidation:
def test_generate_stream_joins_worker_after_response_exhaustion(self, monkeypatch: pytest.MonkeyPatch):
generator = WorkflowAppGenerator()
worker_thread = Mock()
worker_thread.is_alive.return_value = False
app_config = WorkflowUIBasedAppConfig(
tenant_id="tenant",
app_id="app",
app_mode=AppMode.WORKFLOW,
additional_features=AppAdditionalFeatures(),
variables=[],
workflow_id="workflow-id",
)
application_generate_entity = WorkflowAppGenerateEntity.model_construct(
task_id="task",
app_config=app_config,
inputs={},
files=[],
user_id="user",
stream=True,
invoke_from=InvokeFrom.WEB_APP,
extras={},
)
def response_stream():
yield {"event": "workflow_finished"}
monkeypatch.setattr(generator, "_bind_file_access_scope", lambda **kwargs: contextlib.nullcontext())
monkeypatch.setattr(
"core.app.apps.workflow.app_generator.WorkflowAppQueueManager",
lambda **kwargs: SimpleNamespace(**kwargs),
)
monkeypatch.setattr(
"core.app.apps.workflow.app_generator.current_app",
SimpleNamespace(_get_current_object=lambda: SimpleNamespace(name="flask")),
)
monkeypatch.setattr("core.app.apps.workflow.app_generator.contextvars.copy_context", lambda: "ctx")
monkeypatch.setattr("core.app.apps.workflow.app_generator.threading.Thread", lambda **kwargs: worker_thread)
monkeypatch.setattr(
"core.app.apps.workflow.app_generator.db",
SimpleNamespace(session=SimpleNamespace(close=Mock())),
)
monkeypatch.setattr(generator, "_get_draft_var_saver_factory", lambda *args, **kwargs: "draft-factory")
monkeypatch.setattr(generator, "_handle_response", lambda **kwargs: response_stream())
monkeypatch.setattr(
"core.app.apps.workflow.app_generator.WorkflowAppGenerateResponseConverter.convert",
lambda response, invoke_from: response,
)
managed_stream = generator._generate(
app_model=SimpleNamespace(mode=AppMode.WORKFLOW, tenant_id="tenant"),
workflow=SimpleNamespace(id="workflow-id"),
user=SimpleNamespace(id="user"),
application_generate_entity=application_generate_entity,
invoke_from=InvokeFrom.WEB_APP,
workflow_execution_repository=SimpleNamespace(),
workflow_node_execution_repository=SimpleNamespace(),
streaming=True,
)
worker_thread.start.assert_called_once_with()
worker_thread.join.assert_not_called()
assert list(managed_stream) == [{"event": "workflow_finished"}]
worker_thread.join.assert_called_once_with(timeout=300)
def test_ensure_snippet_start_node_returns_original_for_non_snippet_workflow(self):
workflow = SimpleNamespace(kind_or_standard="workflow")
session = SimpleNamespace(scalar=Mock())
@@ -1,67 +0,0 @@
"""Verify that Jinja2 transform_caller scripts block unsafe template attacks via SandboxedEnvironment."""
import io
import sys
import pytest
from core.helper.code_executor.jinja2.jinja2_transformer import Jinja2TemplateTransformer
MALICIOUS_TEMPLATES = [
pytest.param(
"{{ lipsum.__globals__.__builtins__.__import__('os').popen('id').read() }}",
id="lipsum_globals_builtins",
),
pytest.param(
"{{ ''.__class__.__mro__[1].__subclasses__() }}",
id="string_class_mro",
),
pytest.param(
"{{ cycler.__init__.__globals__.os.popen('whoami').read() }}",
id="cycler_init_globals",
),
pytest.param(
"{{ namespace.__init__.__globals__['__builtins__']['__import__']('os').system('id') }}",
id="namespace_init_globals",
),
]
def _exec_scripts(runner: str, preload: str) -> str:
"""Execute preload then runner in a shared namespace, return captured stdout."""
ns: dict = {}
exec(compile(preload, "<preload>", "exec"), ns) # noqa: S102
captured = io.StringIO()
old_stdout = sys.stdout
sys.stdout = captured
try:
exec(compile(runner, "<runner>", "exec"), ns) # noqa: S102
finally:
sys.stdout = old_stdout
return captured.getvalue()
class TestJinja2TransformCallerSandbox:
"""Test transform_caller output (runner + preload) blocks attacks and allows safe templates."""
@pytest.mark.parametrize("malicious_template", MALICIOUS_TEMPLATES)
def test_blocks_unsafe_template(self, malicious_template: str) -> None:
runner, preload = Jinja2TemplateTransformer.transform_caller(malicious_template, {})
ns: dict = {}
exec(compile(preload, "<preload>", "exec"), ns) # noqa: S102
with pytest.raises(Exception) as exc_info:
exec(compile(runner, "<runner>", "exec"), ns) # noqa: S102
assert "unsafe" in str(exc_info.value).lower() or "security" in str(exc_info.value).lower()
def test_renders_safe_template(self) -> None:
runner, preload = Jinja2TemplateTransformer.transform_caller(
"Hello {{ name }}, you are {{ age }} years old!",
{"name": "Alice", "age": 30},
)
output = _exec_scripts(runner, preload)
assert "Hello Alice, you are 30 years old!" in output
def test_scripts_use_sandboxed_environment(self) -> None:
runner, preload = Jinja2TemplateTransformer.transform_caller("{{ x }}", {"x": 1})
assert "SandboxedEnvironment" in runner
assert "SandboxedEnvironment" in preload
@@ -1,57 +1,7 @@
from datetime import datetime, timedelta
import sys
from unittest.mock import MagicMock, patch
import pytest
from sqlalchemy import event
from sqlalchemy.exc import SQLAlchemyError
from sqlalchemy.orm import Session, sessionmaker
import core.llm_generator.llm_generator as generator_module
from core.llm_generator.llm_generator import LLMGenerator, _parse_string_list
from core.model_manager import ModelInstance, ModelManager
from core.workflow.generator import tool_catalogue as tool_catalogue_module
from core.workflow.generator.tool_catalogue import ToolCatalogueEntry
from graphon.model_runtime.entities.llm_entities import LLMResult, LLMUsage
from graphon.model_runtime.entities.message_entities import AssistantPromptMessage
from models.dataset import Dataset
from services.workflow_service import WorkflowService
@pytest.fixture
def dataset_session(sqlite_session: Session, monkeypatch: pytest.MonkeyPatch) -> Session:
"""Bind the real SQLite session to the production database extension."""
monkeypatch.setattr(generator_module.db, "session", sqlite_session)
return sqlite_session
def _llm_result(content: str) -> LLMResult:
"""Build a real non-streaming LLM response around deterministic test content."""
return LLMResult(
model="test-model",
message=AssistantPromptMessage(content=content),
usage=LLMUsage.empty_usage(),
)
def _model_manager() -> tuple[MagicMock, MagicMock]:
"""Build spec-constrained mocks for the model-manager boundary and its default model."""
model_manager = MagicMock(spec=ModelManager)
model_instance = MagicMock(spec=ModelInstance)
model_manager.get_default_model_instance.return_value = model_instance
return model_manager, model_instance
def _dataset(*, dataset_id: str, tenant_id: str, name: str, created_at: datetime) -> Dataset:
return Dataset(
id=dataset_id,
tenant_id=tenant_id,
name=name,
created_by="account-id",
created_at=created_at,
)
class TestParseStringList:
@@ -84,115 +34,95 @@ class TestParseStringList:
class TestGenerateWorkflowInstructionSuggestions:
@patch("core.llm_generator.llm_generator.ModelManager.for_tenant")
def test_no_default_model(self, mock_for_tenant):
model_manager, _ = _model_manager()
model_manager.get_default_model_instance.side_effect = RuntimeError("no default model")
mock_for_tenant.return_value = model_manager
mock_for_tenant.return_value.get_default_model_instance.side_effect = Exception("No model")
assert LLMGenerator.generate_workflow_instruction_suggestions("tenant", mode="workflow") == []
@patch("core.llm_generator.llm_generator.ModelManager.for_tenant")
@patch("core.llm_generator.llm_generator.LLMGenerator._build_suggestion_context")
def test_llm_success(self, mock_build_context, mock_for_tenant):
mock_build_context.return_value = "context"
model_manager, model_instance = _model_manager()
model_instance.invoke_llm.return_value = _llm_result('["idea 1", "idea 2"]')
mock_for_tenant.return_value = model_manager
mock_model = MagicMock()
mock_model.invoke_llm.return_value = MagicMock()
mock_model.invoke_llm.return_value.message.get_text_content.return_value = '["idea 1", "idea 2"]'
mock_for_tenant.return_value.get_default_model_instance.return_value = mock_model
result = LLMGenerator.generate_workflow_instruction_suggestions("tenant", mode="workflow")
assert result == ["idea 1", "idea 2"]
model_instance.invoke_llm.assert_called_once()
@patch("core.llm_generator.llm_generator.ModelManager.for_tenant")
@patch("core.llm_generator.llm_generator.LLMGenerator._build_suggestion_context")
def test_llm_error(self, mock_build_context, mock_for_tenant):
mock_build_context.return_value = "context"
model_manager, model_instance = _model_manager()
model_instance.invoke_llm.side_effect = RuntimeError("API error")
mock_for_tenant.return_value = model_manager
result = LLMGenerator.generate_workflow_instruction_suggestions("tenant", mode="workflow")
assert result == []
model_instance.invoke_llm.assert_called_once()
mock_model = MagicMock()
mock_model.invoke_llm.side_effect = Exception("API error")
mock_for_tenant.return_value.get_default_model_instance.return_value = mock_model
assert LLMGenerator.generate_workflow_instruction_suggestions("tenant", mode="workflow") == []
@patch("core.llm_generator.llm_generator.ModelManager.for_tenant")
@patch("core.llm_generator.llm_generator.LLMGenerator._build_suggestion_context")
def test_llm_bad_output(self, mock_build_context, mock_for_tenant):
mock_build_context.return_value = "context"
model_manager, model_instance = _model_manager()
model_instance.invoke_llm.return_value = _llm_result("Not a list")
mock_for_tenant.return_value = model_manager
result = LLMGenerator.generate_workflow_instruction_suggestions("tenant", mode="workflow")
assert result == []
model_instance.invoke_llm.assert_called_once()
mock_model = MagicMock()
mock_model.invoke_llm.return_value = MagicMock()
mock_model.invoke_llm.return_value.message.get_text_content.return_value = "Not a list"
mock_for_tenant.return_value.get_default_model_instance.return_value = mock_model
assert LLMGenerator.generate_workflow_instruction_suggestions("tenant", mode="workflow") == []
@pytest.mark.parametrize("sqlite_session", [(Dataset,)], indirect=True)
class TestBuildSuggestionContext:
def test_both_success(self, dataset_session: Session, monkeypatch: pytest.MonkeyPatch):
now = datetime.now()
dataset_session.add_all(
(
_dataset(dataset_id="kb-1", tenant_id="tenant", name="kb1", created_at=now),
_dataset(
dataset_id="kb-2",
tenant_id="tenant",
name="kb2",
created_at=now - timedelta(seconds=1),
),
_dataset(dataset_id="other-kb", tenant_id="other", name="private", created_at=now),
)
)
dataset_session.commit()
@patch("core.llm_generator.llm_generator.db.session.scalars")
def test_both_success(self, mock_scalars, monkeypatch):
mock_scalars.return_value.all.return_value = ["kb1", "kb2"]
def build_tool_catalogue(_tenant_id: str) -> list[ToolCatalogueEntry]:
return [
ToolCatalogueEntry(
provider_name="provider",
provider_type="builtin",
plugin_id="",
tool_name="tool1",
tool_label="tool1",
description="First tool",
),
ToolCatalogueEntry(
provider_name="provider",
provider_type="builtin",
plugin_id="",
tool_name="tool2",
tool_label="tool2",
description="Second tool",
),
]
# Keep the real module and formatter; only isolate provider/plugin discovery.
monkeypatch.setattr(tool_catalogue_module, "build_tool_catalogue", build_tool_catalogue)
# ``_build_suggestion_context`` imports the tool catalogue lazily, so we
# stub the module in ``sys.modules``. Use ``monkeypatch.setitem`` so the
# ORIGINAL module is RESTORED on teardown — a bare ``del`` would evict it
# from sys.modules entirely, after which a sibling test that imported
# ``build_tool_catalogue`` at collection time (e.g. test_tool_catalogue)
# diverges from a freshly re-imported module and its @patch targets stop
# applying, silently breaking it under xdist.
mock_tool_catalogue = MagicMock()
mock_tool_catalogue.build_tool_catalogue.return_value = "catalog"
mock_tool_catalogue.format_tool_catalogue.return_value = "tool1\ntool2"
monkeypatch.setitem(sys.modules, "core.workflow.generator.tool_catalogue", mock_tool_catalogue)
result = LLMGenerator._build_suggestion_context("tenant")
assert "Knowledge bases:\n- kb1\n- kb2" in result
assert "Installed tools:\n- provider/tool1 — First tool\n- provider/tool2 — Second tool" in result
assert "Installed tools:\ntool1\ntool2" in result
def test_both_fail(self, dataset_session: Session, monkeypatch: pytest.MonkeyPatch):
def fail_query(_orm_execute_state: object) -> None:
raise SQLAlchemyError("DB error")
@patch("core.llm_generator.llm_generator.db.session.scalars")
def test_both_fail(self, mock_scalars, monkeypatch):
mock_scalars.side_effect = Exception("DB error")
def fail_tool_catalogue(_tenant_id: str) -> list[ToolCatalogueEntry]:
raise RuntimeError("Tool error")
# See ``test_both_success``: restore the original module via monkeypatch
# rather than ``del``-ing it, so we don't evict it for sibling tests.
mock_tool_catalogue = MagicMock()
mock_tool_catalogue.build_tool_catalogue.side_effect = Exception("Tool error")
monkeypatch.setitem(sys.modules, "core.workflow.generator.tool_catalogue", mock_tool_catalogue)
event.listen(dataset_session, "do_orm_execute", fail_query)
monkeypatch.setattr(tool_catalogue_module, "build_tool_catalogue", fail_tool_catalogue)
try:
assert LLMGenerator._build_suggestion_context("tenant") == ""
finally:
event.remove(dataset_session, "do_orm_execute", fail_query)
assert LLMGenerator._build_suggestion_context("tenant") == ""
class TestWorkflowServiceInterface:
def test_real_workflow_service_exposes_protocol_methods(self):
def test_protocol_methods(self):
# Just to cover the 'pass' statements in the Protocol definition
from core.llm_generator.llm_generator import WorkflowServiceInterface
service: WorkflowServiceInterface = WorkflowService(sessionmaker())
class MockService(WorkflowServiceInterface):
def get_draft_workflow(self, app_model, workflow_id=None, *, session):
return super().get_draft_workflow(app_model, workflow_id, session=session)
assert callable(service.get_draft_workflow)
assert callable(service.get_node_last_run)
def get_node_last_run(self, app_model, workflow, node_id):
return super().get_node_last_run(app_model, workflow, node_id)
service = MockService()
service.get_draft_workflow(None, session=None)
service.get_node_last_run(None, None, "node")
@@ -2,24 +2,13 @@ from unittest.mock import MagicMock, patch
import pytest
from pydantic import ValidationError
from sqlalchemy.orm import Session
import core.moderation.api.api as moderation_module
from core.extension.api_based_extension_requestor import APIBasedExtensionPoint
from core.moderation.api.api import ApiModeration, ModerationInputParams, ModerationOutputParams
from core.moderation.base import ModerationAction, ModerationInputsResult, ModerationOutputsResult
from models.api_based_extension import APIBasedExtension
class _DatabaseBinding:
"""Expose the real SQLite session used by extension lookup."""
session: Session
def __init__(self, session: Session) -> None:
self.session = session
class TestApiModeration:
@pytest.fixture
def api_config(self):
@@ -176,27 +165,17 @@ class TestApiModeration:
with pytest.raises(ValueError, match="API-based Extension not found"):
api_moderation._get_config_by_requestor(APIBasedExtensionPoint.APP_MODERATION_INPUT, {})
@pytest.mark.parametrize("sqlite_session", [(APIBasedExtension,)], indirect=True)
def test_get_api_based_extension(self, sqlite_session: Session, monkeypatch: pytest.MonkeyPatch) -> None:
target = APIBasedExtension(
tenant_id="tenant-1",
name="Target extension",
api_endpoint="https://example.com/moderate",
api_key="encrypted-key",
)
target.id = "ext-1"
other_tenant = APIBasedExtension(
tenant_id="tenant-2",
name="Other extension",
api_endpoint="https://example.com/other",
api_key="other-key",
)
other_tenant.id = "ext-2"
sqlite_session.add_all((target, other_tenant))
sqlite_session.commit()
monkeypatch.setattr(moderation_module, "db", _DatabaseBinding(sqlite_session))
@patch("core.moderation.api.api.db.session.scalar")
def test_get_api_based_extension(self, mock_scalar):
mock_ext = MagicMock(spec=APIBasedExtension)
mock_scalar.return_value = mock_ext
result = ApiModeration._get_api_based_extension("tenant-1", "ext-1")
assert result is target
assert ApiModeration._get_api_based_extension("tenant-1", "ext-2") is None
assert result == mock_ext
mock_scalar.assert_called_once()
# Verify the call has the correct filters
args, kwargs = mock_scalar.call_args
stmt = args[0]
# We can't easily inspect the statement without complex sqlalchemy tricks,
# but calling it is usually enough for unit tests if we mock the result.
+11 -47
View File
@@ -1,11 +1,9 @@
import re
from datetime import datetime
from decimal import Decimal
from unittest.mock import MagicMock, patch
import pytest
from sqlalchemy.orm import Session
import core.ops.utils as utils_module
from core.ops.utils import (
filter_none_values,
generate_dotted_order,
@@ -17,42 +15,6 @@ from core.ops.utils import (
validate_url,
validate_url_with_path,
)
from models.enums import ConversationFromSource
from models.model import Message
class _DatabaseBinding:
"""Expose the real SQLite session used by the message lookup helper."""
session: Session
def __init__(self, session: Session) -> None:
self.session = session
@pytest.fixture
def message_session(sqlite_session: Session, monkeypatch: pytest.MonkeyPatch) -> Session:
"""Bind the message lookup helper to the shared SQLite test session."""
monkeypatch.setattr(utils_module, "db", _DatabaseBinding(sqlite_session))
return sqlite_session
def _message(message_id: str) -> Message:
message = Message(
id=message_id,
app_id="app-id",
conversation_id="conversation-id",
query="question",
message={"role": "user", "content": "question"},
answer="answer",
message_unit_price=Decimal("0.0001"),
answer_unit_price=Decimal("0.0001"),
currency="USD",
from_source=ConversationFromSource.API,
)
message._inputs = {}
return message
class TestValidateUrl:
@@ -258,20 +220,22 @@ class TestFilterNoneValues:
assert filter_none_values({}) == {}
@pytest.mark.parametrize("sqlite_session", [(Message,)], indirect=True)
class TestGetMessageData:
"""Test cases for get_message_data function"""
def test_get_message_data(self, message_session: Session):
target = _message("message-id")
unrelated = _message("other-message-id")
message_session.add_all((target, unrelated))
message_session.commit()
@patch("core.ops.utils.db")
@patch("core.ops.utils.Message")
@patch("core.ops.utils.select")
def test_get_message_data(self, mock_select, mock_message, mock_db):
mock_scalar = mock_db.session.scalar
mock_msg_instance = MagicMock()
mock_scalar.return_value = mock_msg_instance
result = get_message_data("message-id")
assert result is target
assert result.id == "message-id"
assert result == mock_msg_instance
mock_select.assert_called_once()
mock_scalar.assert_called_once()
class TestMeasureTime:
@@ -1,42 +1,9 @@
from datetime import datetime, timedelta
from decimal import Decimal
from unittest.mock import MagicMock
from uuid import uuid4
import pytest
from sqlalchemy.orm import Session
from constants import UUID_NIL
from core.prompt.utils.extract_thread_messages import extract_thread_messages
from core.prompt.utils.get_thread_messages_length import get_thread_messages_length
from models.enums import ConversationFromSource
from models.model import Message
def _persisted_message(
*,
message_id: str,
conversation_id: str,
parent_message_id: str,
answer: str,
created_at: datetime,
) -> Message:
message = Message(
id=message_id,
app_id="app-id",
conversation_id=conversation_id,
query="question",
message={"role": "user", "content": "question"},
answer=answer,
message_unit_price=Decimal("0.0001"),
answer_unit_price=Decimal("0.0001"),
currency="USD",
from_source=ConversationFromSource.API,
parent_message_id=parent_message_id,
created_at=created_at,
updated_at=created_at,
)
message._inputs = {}
return message
class MockMessage:
@@ -137,64 +104,33 @@ def test_extract_thread_messages_breaks_when_parent_is_none():
assert result[0].id == id2
@pytest.mark.parametrize("sqlite_session", [(Message,)], indirect=True)
def test_get_thread_messages_length_excludes_newly_created_empty_answer(sqlite_session: Session):
def test_get_thread_messages_length_excludes_newly_created_empty_answer():
id1, id2 = str(uuid4()), str(uuid4())
now = datetime.now()
messages = [
_persisted_message(
message_id=id2,
conversation_id="conversation-1",
parent_message_id=id1,
answer="",
created_at=now,
),
_persisted_message(
message_id=id1,
conversation_id="conversation-1",
parent_message_id=UUID_NIL,
answer="ok",
created_at=now - timedelta(seconds=1),
),
_persisted_message(
message_id=str(uuid4()),
conversation_id="other-conversation",
parent_message_id=UUID_NIL,
answer="unrelated",
created_at=now + timedelta(seconds=1),
),
MockMessage(id2, id1, answer=""), # newest generated message should be excluded
MockMessage(id1, UUID_NIL, answer="ok"),
]
sqlite_session.add_all(messages)
sqlite_session.commit()
length = get_thread_messages_length("conversation-1", session=sqlite_session)
session = MagicMock()
session.scalars.return_value.all.return_value = messages
length = get_thread_messages_length("conversation-1", session=session)
assert length == 1
session.scalars.assert_called_once()
@pytest.mark.parametrize("sqlite_session", [(Message,)], indirect=True)
def test_get_thread_messages_length_keeps_non_empty_latest_answer(sqlite_session: Session):
def test_get_thread_messages_length_keeps_non_empty_latest_answer():
id1, id2 = str(uuid4()), str(uuid4())
now = datetime.now()
messages = [
_persisted_message(
message_id=id2,
conversation_id="conversation-2",
parent_message_id=id1,
answer="latest-answer",
created_at=now,
),
_persisted_message(
message_id=id1,
conversation_id="conversation-2",
parent_message_id=UUID_NIL,
answer="older-answer",
created_at=now - timedelta(seconds=1),
),
MockMessage(id2, id1, answer="latest-answer"),
MockMessage(id1, UUID_NIL, answer="older-answer"),
]
sqlite_session.add_all(messages)
sqlite_session.commit()
length = get_thread_messages_length("conversation-2", session=sqlite_session)
session = MagicMock()
session.scalars.return_value.all.return_value = messages
length = get_thread_messages_length("conversation-2", session=session)
assert length == 2
session.scalars.assert_called_once()
@@ -22,23 +22,6 @@ class TestCleanProcessor:
expected = "normalpadding"
assert CleanProcessor.clean(text_with_ufffe, None) == expected
def test_clean_preserves_valid_extended_characters(self):
"""Default cleaning must not strip valid printable characters.
The invalid-symbol filter used to include the UTF-8 bytes of U+FFFE
(0xEF 0xBF 0xBE) inside a character class. On a decoded string those
bytes are the code points U+00EF, U+00BF and U+00BE, i.e. the valid
characters 'ï', '¿' and '¾', so words like "naïve" and Spanish
questions like "¿Cómo?" were being silently corrupted on ingest.
"""
assert CleanProcessor.clean("naïve", None) == "naïve"
assert CleanProcessor.clean("¿Cómo estás?", None) == "¿Cómo estás?"
assert CleanProcessor.clean("¾ cup sugar", None) == "¾ cup sugar"
assert CleanProcessor.clean("￾", None) == "￾"
# The U+FFFE noncharacter is still stripped by its dedicated substitution.
assert CleanProcessor.clean("keep\ufffedrop", None) == "keepdrop"
def test_clean_with_none_process_rule(self):
"""Test cleaning with None process_rule - only default cleaning applied."""
text = "Hello<|World\x00"
@@ -1372,19 +1372,6 @@ class TestIndexingRunnerDocumentCleaning:
assert "\ufffe" not in result
assert "Text with" in result
def test_filter_string_preserves_valid_extended_characters(self):
"""filter_string must keep valid printable characters like 'ï', '¿', '¾'."""
# Arrange
text = "naïve ¿Cómo? ¾ done"
# Act
result = IndexingRunner.filter_string(text)
# Assert
assert result == text
# The U+FFFE noncharacter is still stripped.
assert IndexingRunner.filter_string("keep\ufffedrop") == "keepdrop"
class TestIndexingRunnerSplitter:
"""Unit tests for text splitter configuration.
@@ -952,21 +952,6 @@ class TestFixedRecursiveCharacterTextSplitter:
assert "word1" in combined
assert "word2" in combined
def test_preserves_spaces_when_recursively_splitting_long_paragraph(self):
"""Ensure recursive space splitting preserves word boundaries."""
text = "여름철에는 항상 기상상황에 주목하며 주변 사람들과 함께 정보를 공유합니다."
splitter = FixedRecursiveCharacterTextSplitter(
fixed_separator="\n\n",
chunk_size=20,
chunk_overlap=0,
keep_separator=True,
)
result = splitter.split_text(text)
assert len(result) > 1
assert " ".join(result) == text
def test_character_level_splitting(self):
"""Test character-level splitting when no separator works."""
text = "verylongwordwithoutspaces"
@@ -18,6 +18,7 @@ from models.provider import ProviderType
@pytest.fixture
def credit_pool_session_factory(sqlite_engine: Engine) -> Iterator[sessionmaker[Session]]:
"""Bind message-created accounting to fixture-owned SQLite sessions."""
TenantCreditPool.__table__.create(sqlite_engine)
session_factory = sessionmaker(bind=sqlite_engine, expire_on_commit=False)
with patch("events.event_handlers.update_provider_when_message_created.db.session", session_factory):
yield session_factory
@@ -1,43 +0,0 @@
from types import SimpleNamespace
from unittest.mock import Mock
from extensions.ext_celery import _enqueue_initial_community_telemetry_heartbeat
def test_beat_start_enqueues_community_telemetry_heartbeat() -> None:
task = Mock()
sender = SimpleNamespace(
app=SimpleNamespace(
conf=SimpleNamespace(beat_schedule={"community_telemetry_heartbeat": {}}),
tasks={"community_telemetry.send_heartbeat": task},
)
)
_enqueue_initial_community_telemetry_heartbeat(sender)
task.apply_async.assert_called_once_with()
def test_beat_start_skips_community_telemetry_when_not_scheduled() -> None:
task = Mock()
sender = SimpleNamespace(
app=SimpleNamespace(
conf=SimpleNamespace(beat_schedule={}),
tasks={"community_telemetry.send_heartbeat": task},
)
)
_enqueue_initial_community_telemetry_heartbeat(sender)
task.apply_async.assert_not_called()
def test_beat_start_skips_community_telemetry_when_task_is_unavailable() -> None:
sender = SimpleNamespace(
app=SimpleNamespace(
conf=SimpleNamespace(beat_schedule={"community_telemetry_heartbeat": {}}),
tasks={},
)
)
_enqueue_initial_community_telemetry_heartbeat(sender)
@@ -548,7 +548,6 @@ def test_load_agent_app_composer_exposes_draft_save_only(monkeypatch: pytest.Mon
agent = SimpleNamespace(
id="agent-1",
active_config_snapshot_id="version-1",
active_config_is_published=True,
updated_by="account-1",
created_by="account-1",
app_id="app-1",
@@ -568,7 +567,6 @@ def test_load_agent_app_composer_exposes_draft_save_only(monkeypatch: pytest.Mon
result = AgentComposerService.load_agent_app_composer(session=session, tenant_id="tenant-1", app_id="app-1")
assert result["save_options"] == [ComposerSaveStrategy.SAVE_TO_CURRENT_VERSION.value]
assert result["active_config_is_published"] is True
def test_save_agent_app_composer_rejects_version_save_strategy():
@@ -612,11 +610,7 @@ def test_save_agent_app_composer_updates_normal_draft(monkeypatch: pytest.Monkey
lambda **kwargs: saved.update(kwargs) or SimpleNamespace(id="draft-1"),
)
monkeypatch.setattr(AgentComposerService, "_get_version_if_present", lambda **_kwargs: active_version)
monkeypatch.setattr(
AgentComposerService,
"load_agent_composer",
lambda **kwargs: {"loaded": True, "active_config_is_published": agent.active_config_is_published},
)
monkeypatch.setattr(AgentComposerService, "load_agent_composer", lambda **kwargs: {"loaded": True})
payload = ComposerSavePayload.model_validate(
{
"variant": ComposerVariant.AGENT_APP.value,
@@ -634,7 +628,7 @@ def test_save_agent_app_composer_updates_normal_draft(monkeypatch: pytest.Monkey
)
assert result.pop("validation") == {"warnings": [], "knowledge_retrieval_placeholder": []}
assert result == {"loaded": True, "active_config_is_published": False}
assert result == {"loaded": True}
assert saved["draft_type"] == AgentConfigDraftType.DRAFT
assert saved["agent_soul"].model_dump(mode="json") == _agent_soul_with_model().model_dump(mode="json")
assert agent.active_config_is_published is False
@@ -663,11 +657,7 @@ def test_save_agent_app_composer_keeps_published_when_draft_matches_active_snaps
lambda **_kwargs: SimpleNamespace(id="draft-1"),
)
monkeypatch.setattr(AgentComposerService, "_get_version_if_present", lambda **_kwargs: active_version)
monkeypatch.setattr(
AgentComposerService,
"load_agent_composer",
lambda **_kwargs: {"loaded": True, "active_config_is_published": agent.active_config_is_published},
)
monkeypatch.setattr(AgentComposerService, "load_agent_composer", lambda **_kwargs: {"loaded": True})
payload = ComposerSavePayload.model_validate(
{
"variant": ComposerVariant.AGENT_APP.value,
@@ -676,7 +666,7 @@ def test_save_agent_app_composer_keeps_published_when_draft_matches_active_snaps
}
)
result = AgentComposerService.save_agent_app_composer(
AgentComposerService.save_agent_app_composer(
session=session,
tenant_id="tenant-1",
app_id="app-1",
@@ -685,7 +675,6 @@ def test_save_agent_app_composer_keeps_published_when_draft_matches_active_snaps
)
assert agent.active_config_is_published is True
assert result["active_config_is_published"] is True
assert fake_session.flushes >= 1
@@ -246,26 +246,6 @@ class TestPluginModelProviderCache:
call([cache_key]),
]
def test_fetch_plugin_model_providers_bypasses_redis_when_cache_disabled(self) -> None:
"""With the cache disabled the daemon is the only source, and Redis is never touched."""
with patch(f"{MODULE}.redis_client") as redis_client, patch(f"{MODULE}.dify_config") as config:
config.PLUGIN_MODEL_PROVIDERS_CACHE_ENABLED = False
client = Mock()
client.fetch_model_providers.return_value = [_build_plugin_model_provider()]
from core.plugin.plugin_service import PluginService
first = PluginService.fetch_plugin_model_providers(tenant_id="tenant-1", client=client)
second = PluginService.fetch_plugin_model_providers(tenant_id="tenant-1", client=client)
assert [provider.provider for provider in first] == ["langgenius/openai/openai"]
assert [provider.provider for provider in second] == ["langgenius/openai/openai"]
assert client.fetch_model_providers.call_count == 2
redis_client.get.assert_not_called()
redis_client.mget.assert_not_called()
redis_client.setex.assert_not_called()
redis_client.lock.assert_not_called()
def test_fetch_plugin_model_providers_refetches_when_cache_read_fails(self) -> None:
"""Redis read failures do not block provider discovery for the tenant."""
with patch(f"{MODULE}.redis_client") as redis_client:
@@ -1,7 +1,7 @@
import json
from collections.abc import Iterator
from datetime import datetime, timedelta
from unittest.mock import MagicMock, patch
from uuid import UUID
import pytest
from sqlalchemy import event, select
@@ -10,6 +10,7 @@ from sqlalchemy.orm import Session
from configs import dify_config
from models.account import (
Account,
AccountIntegrate,
AccountStatus,
Tenant,
TenantAccountJoin,
@@ -112,6 +113,22 @@ class TestAccountService:
- Error conditions and edge cases
"""
@pytest.fixture
def sqlite_session(self, sqlite_engine) -> Iterator[Session]:
"""SQLite session with the account/workspace tables these service tests touch."""
tables = [
model.metadata.tables[model.__tablename__]
for model in (
Account,
Tenant,
TenantAccountJoin,
TenantPluginAutoUpgradeStrategy,
)
]
Account.metadata.create_all(sqlite_engine, tables=tables)
with Session(sqlite_engine, expire_on_commit=False) as session:
yield session
@pytest.fixture
def mock_password_dependencies(self):
"""Mock setup for password-related functions."""
@@ -1246,6 +1263,24 @@ class TestRegisterService:
- Error conditions and edge cases
"""
@pytest.fixture
def sqlite_session(self, sqlite_engine) -> Iterator[Session]:
"""SQLite session with the account/workspace tables registration flows touch."""
tables = [
model.metadata.tables[model.__tablename__]
for model in (
Account,
AccountIntegrate,
Tenant,
TenantAccountJoin,
TenantPluginAutoUpgradeStrategy,
DifySetup,
)
]
Account.metadata.create_all(sqlite_engine, tables=tables)
with Session(sqlite_engine, expire_on_commit=False) as session:
yield session
@pytest.fixture
def mock_redis_dependencies(self):
"""Mock setup for Redis-related functions."""
@@ -1290,10 +1325,7 @@ class TestRegisterService:
with patch("services.account_service.AccountService.create_account") as mock_create_account:
mock_create_account.return_value = mock_account
with (
patch("services.account_service.TenantService.create_owner_tenant_if_not_exist") as mock_create_tenant,
patch("services.account_service.CommunityTelemetryService.report_install") as mock_report_install,
):
with patch("services.account_service.TenantService.create_owner_tenant_if_not_exist") as mock_create_tenant:
RegisterService.setup(
"[email protected]",
"Admin User",
@@ -1312,39 +1344,7 @@ class TestRegisterService:
session=sqlite_session,
)
mock_create_tenant.assert_called_once_with(account=mock_account, is_setup=True, session=sqlite_session)
dify_setup = sqlite_session.scalar(select(DifySetup))
assert dify_setup is not None
assert dify_setup.instance_id is not None
assert str(UUID(dify_setup.instance_id)) == dify_setup.instance_id
assert dify_setup.install_reported_at is None
assert dify_setup.last_heartbeat_at is None
mock_report_install.assert_called_once_with(session=sqlite_session)
def test_setup_succeeds_when_telemetry_install_report_fails(
self, sqlite_session: Session, mock_external_service_dependencies
):
mock_external_service_dependencies["feature_service"].get_system_features.return_value.is_allow_register = True
mock_external_service_dependencies["billing_service"].is_email_in_freeze.return_value = False
mock_account = TestAccountAssociatedDataFactory.create_account_mock()
with (
patch("services.account_service.AccountService.create_account", return_value=mock_account),
patch("services.account_service.TenantService.create_owner_tenant_if_not_exist"),
patch(
"services.account_service.CommunityTelemetryService.report_install",
side_effect=RuntimeError("telemetry unavailable"),
),
):
RegisterService.setup(
"[email protected]",
"Admin User",
"password123",
"192.168.1.1",
"en-US",
session=sqlite_session,
)
assert sqlite_session.scalar(select(DifySetup)) is not None
assert sqlite_session.scalar(select(DifySetup)) is not None
def test_setup_failure_rollback(self, sqlite_session: Session, mock_external_service_dependencies):
"""Test setup failure with proper rollback."""
@@ -1,23 +1,19 @@
from __future__ import annotations
from collections.abc import Callable
from datetime import datetime
from types import SimpleNamespace
from typing import cast
from unittest.mock import MagicMock, patch
from uuid import uuid4
import pytest
from sqlalchemy import event
from sqlalchemy.exc import IntegrityError
from sqlalchemy.orm import Session
from graphon.model_runtime.entities.model_entities import ModelType
from models import Account
from models.model import App, AppMode, AppModelConfig, IconType
from models.model import App, AppMode, AppModelConfig
from models.workflow import Workflow
from services.agent.errors import AgentNameConflictError
from services.app_service import AppListParams, AppService, CreateAppParams
from services.app_service import AppService, CreateAppParams
class TestCreateAppTransactionBoundary:
@@ -240,92 +236,6 @@ class TestOpenapiVisibilityHelpers:
mock_session.execute.assert_called_once()
@pytest.mark.parametrize("sqlite_session", [(Account, App, AppModelConfig)], indirect=True)
def test_get_recent_apps_uses_one_tenant_scoped_projection_query(sqlite_session: Session) -> None:
tenant_id = str(uuid4())
other_tenant_id = str(uuid4())
account = Account(name="Recent Apps Author", email="[email protected]")
sqlite_session.add(account)
sqlite_session.flush()
def create_app(*, name: str, tenant_id: str, updated_at: datetime, mode: AppMode = AppMode.CHAT) -> App:
app = App()
app.id = str(uuid4())
app.tenant_id = tenant_id
app.name = name
app.description = ""
app.mode = mode
app.icon_type = IconType.EMOJI
app.icon = "🚀"
app.icon_background = "#FFFFFF"
app.enable_site = False
app.enable_api = False
app.created_by = account.id
app.maintainer = account.id
app.created_at = updated_at
app.updated_at = updated_at
app.use_icon_as_answer_icon = False
return app
newest = create_app(name="Newest", tenant_id=tenant_id, updated_at=datetime(2026, 7, 3))
legacy_agent = AppModelConfig(app_id=newest.id)
legacy_agent.agent_mode = '{"enabled": true, "strategy": "react"}'
newest.app_model_config_id = legacy_agent.id
second = create_app(
name="Second",
tenant_id=tenant_id,
updated_at=datetime(2026, 7, 2),
mode=AppMode.WORKFLOW,
)
second.icon_type = None
second.icon = None
second.icon_background = None
second.created_by = None
second.maintainer = None
channel = create_app(
name="Channel",
tenant_id=tenant_id,
updated_at=datetime(2026, 7, 5),
mode=AppMode.CHANNEL,
)
rag_pipeline = create_app(
name="RAG Pipeline",
tenant_id=tenant_id,
updated_at=datetime(2026, 7, 4),
mode=AppMode.RAG_PIPELINE,
)
oldest = create_app(name="Oldest", tenant_id=tenant_id, updated_at=datetime(2026, 7, 1))
foreign = create_app(name="Foreign", tenant_id=other_tenant_id, updated_at=datetime(2026, 7, 4))
sqlite_session.add_all([newest, legacy_agent, second, channel, rag_pipeline, oldest, foreign])
sqlite_session.commit()
statements: list[str] = []
bind = sqlite_session.get_bind()
def record_sql(_conn, _cursor, statement, _parameters, _context, _executemany) -> None:
statements.append(statement)
event.listen(bind, "before_cursor_execute", record_sql)
try:
recent_apps = AppService().get_recent_apps(
account.id,
tenant_id,
AppListParams(limit=2),
sqlite_session,
)
finally:
event.remove(bind, "before_cursor_execute", record_sql)
assert [(app.name, app.mode, app.icon_type, app.author_name, app.maintainer) for app in recent_apps] == [
("Newest", AppMode.CHAT, IconType.EMOJI, "Recent Apps Author", account.id),
("Second", AppMode.WORKFLOW, None, None, None),
]
select_statements = [statement for statement in statements if statement.lstrip().upper().startswith("SELECT")]
assert len(select_statements) == 1
assert "count(" not in select_statements[0].lower()
assert "app_model_configs" not in select_statements[0].lower()
class TestAppMeta:
def test_loads_workflow_with_caller_session(self):
session = MagicMock()
File diff suppressed because it is too large Load Diff
@@ -1,333 +0,0 @@
import uuid
from datetime import datetime
from unittest.mock import Mock
import httpx
import pytest
from sqlalchemy import select
from sqlalchemy.orm import Session
from models.model import DifySetup
from services import telemetry_service
from services.telemetry_service import CommunityTelemetryService
@pytest.fixture
def telemetry_enabled(monkeypatch: pytest.MonkeyPatch):
monkeypatch.setattr(telemetry_service.dify_config, "EDITION", "SELF_HOSTED")
monkeypatch.setattr(telemetry_service.dify_config, "ENTERPRISE_ENABLED", False)
monkeypatch.setattr(telemetry_service.dify_config, "DISABLE_TELEMETRY", False)
monkeypatch.setattr(telemetry_service.dify_config, "DO_NOT_TRACK", False)
monkeypatch.setattr(telemetry_service.dify_config, "CI", False)
monkeypatch.setattr(telemetry_service.dify_config, "TELEMETRY_ENDPOINT", "https://telemetry.example.test/v1/events")
monkeypatch.setattr(
telemetry_service.dify_config,
"TELEMETRY_FALLBACK_ENDPOINT",
"https://telemetry-cn.example.test/v1/events",
)
monkeypatch.setattr(telemetry_service.dify_config, "TELEMETRY_TIMEOUT_SECONDS", 2)
def test_telemetry_is_disabled_for_enterprise(monkeypatch: pytest.MonkeyPatch):
monkeypatch.setattr(telemetry_service.dify_config, "EDITION", "SELF_HOSTED")
monkeypatch.setattr(telemetry_service.dify_config, "ENTERPRISE_ENABLED", True)
assert CommunityTelemetryService._is_enabled() is False
@pytest.mark.parametrize(
("setting", "value"),
[
("EDITION", "CLOUD"),
("DISABLE_TELEMETRY", True),
("DO_NOT_TRACK", True),
("CI", True),
("TELEMETRY_ENDPOINT", ""),
],
)
def test_telemetry_is_disabled_when_a_required_condition_is_not_met(
telemetry_enabled, monkeypatch: pytest.MonkeyPatch, setting: str, value: str | bool
):
monkeypatch.setattr(telemetry_service.dify_config, setting, value)
assert CommunityTelemetryService._is_enabled() is False
@pytest.mark.parametrize("sqlite_session", [(DifySetup,)], indirect=True)
def test_reporting_without_setup_is_skipped(sqlite_session: Session, telemetry_enabled):
assert CommunityTelemetryService.report_install(session=sqlite_session) is False
assert CommunityTelemetryService.report_heartbeat(session=sqlite_session) is False
@pytest.mark.parametrize("sqlite_session", [(DifySetup,)], indirect=True)
def test_report_install_marks_reported_at(sqlite_session: Session, telemetry_enabled, monkeypatch: pytest.MonkeyPatch):
setup = DifySetup(version="installed-version", instance_id="d246c3a1-350b-406c-92c7-6043df680758")
sqlite_session.add(setup)
sqlite_session.commit()
monkeypatch.setattr(telemetry_service.dify_config.project, "version", "running-version")
sent_payloads: list[dict[str, str | int]] = []
def fake_post(url: str, json: dict[str, str | int], timeout: int):
sent_payloads.append(json)
return httpx.Response(204, request=httpx.Request("POST", url))
monkeypatch.setattr(telemetry_service.httpx, "post", fake_post)
assert CommunityTelemetryService.report_install(session=sqlite_session) is True
saved_setup = sqlite_session.scalar(select(DifySetup))
assert saved_setup is not None
assert saved_setup.install_reported_at is not None
assert sent_payloads[0]["event"] == "install"
assert sent_payloads[0]["instance_id"] == setup.instance_id
assert sent_payloads[0]["version"] == "installed-version"
assert "installed_at" in sent_payloads[0]
@pytest.mark.parametrize("sqlite_session", [(DifySetup,)], indirect=True)
def test_report_install_generates_missing_instance_id(
sqlite_session: Session, telemetry_enabled, monkeypatch: pytest.MonkeyPatch
):
setup = DifySetup(version="installed-version")
sqlite_session.add(setup)
sqlite_session.commit()
monkeypatch.setattr(
telemetry_service.httpx,
"post",
lambda url, json, timeout: httpx.Response(204, request=httpx.Request("POST", url)),
)
assert CommunityTelemetryService.report_install(session=sqlite_session) is True
assert setup.instance_id is not None
assert str(uuid.UUID(setup.instance_id)) == setup.instance_id
@pytest.mark.parametrize("sqlite_session", [(DifySetup,)], indirect=True)
def test_report_heartbeat_generates_missing_instance_id(
sqlite_session: Session, telemetry_enabled, monkeypatch: pytest.MonkeyPatch
):
setup = DifySetup(version="1.0.0", install_reported_at=datetime(2026, 7, 12, 8, 0, 0))
sqlite_session.add(setup)
sqlite_session.commit()
monkeypatch.setattr(
telemetry_service.httpx,
"post",
lambda url, json, timeout: httpx.Response(204, request=httpx.Request("POST", url)),
)
assert (
CommunityTelemetryService.report_heartbeat(session=sqlite_session, now=datetime(2026, 7, 13, 12, 0, 0)) is True
)
assert setup.instance_id is not None
assert str(uuid.UUID(setup.instance_id)) == setup.instance_id
@pytest.mark.parametrize("sqlite_session", [(DifySetup,)], indirect=True)
def test_report_install_failure_keeps_install_pending(
sqlite_session: Session, telemetry_enabled, monkeypatch: pytest.MonkeyPatch
):
setup = DifySetup(version="1.0.0", instance_id="d246c3a1-350b-406c-92c7-6043df680758")
sqlite_session.add(setup)
sqlite_session.commit()
def fake_post(url: str, json: dict[str, str | int], timeout: int):
raise httpx.ConnectError("offline", request=httpx.Request("POST", url))
monkeypatch.setattr(telemetry_service.httpx, "post", fake_post)
assert CommunityTelemetryService.report_install(session=sqlite_session) is False
saved_setup = sqlite_session.scalar(select(DifySetup))
assert saved_setup is not None
assert saved_setup.install_reported_at is None
@pytest.mark.parametrize("sqlite_session", [(DifySetup,)], indirect=True)
def test_report_install_uses_fallback_endpoint_after_network_failure(
sqlite_session: Session, telemetry_enabled, monkeypatch: pytest.MonkeyPatch
):
setup = DifySetup(version="1.0.0", instance_id="d246c3a1-350b-406c-92c7-6043df680758")
sqlite_session.add(setup)
sqlite_session.commit()
urls: list[str] = []
def fake_post(url: str, json: dict[str, str | int], timeout: int):
urls.append(url)
if url == telemetry_service.dify_config.TELEMETRY_ENDPOINT:
raise httpx.ConnectError("offline", request=httpx.Request("POST", url))
return httpx.Response(204, request=httpx.Request("POST", url))
monkeypatch.setattr(telemetry_service.httpx, "post", fake_post)
assert CommunityTelemetryService.report_install(session=sqlite_session) is True
assert urls == [
telemetry_service.dify_config.TELEMETRY_ENDPOINT,
telemetry_service.dify_config.TELEMETRY_FALLBACK_ENDPOINT,
]
@pytest.mark.parametrize("sqlite_session", [(DifySetup,)], indirect=True)
def test_report_install_does_not_use_fallback_endpoint_after_http_error(
sqlite_session: Session, telemetry_enabled, monkeypatch: pytest.MonkeyPatch
):
setup = DifySetup(version="1.0.0", instance_id="d246c3a1-350b-406c-92c7-6043df680758")
sqlite_session.add(setup)
sqlite_session.commit()
post_mock = Mock(
return_value=httpx.Response(
500,
request=httpx.Request("POST", telemetry_service.dify_config.TELEMETRY_ENDPOINT),
)
)
monkeypatch.setattr(telemetry_service.httpx, "post", post_mock)
assert CommunityTelemetryService.report_install(session=sqlite_session) is False
post_mock.assert_called_once()
@pytest.mark.parametrize("sqlite_session", [(DifySetup,)], indirect=True)
def test_report_heartbeat_retries_pending_install_before_heartbeat(
sqlite_session: Session, telemetry_enabled, monkeypatch: pytest.MonkeyPatch
):
setup = DifySetup(version="installed-version", instance_id="d246c3a1-350b-406c-92c7-6043df680758")
sqlite_session.add(setup)
sqlite_session.commit()
monkeypatch.setattr(telemetry_service.dify_config.project, "version", "running-version")
sent_payloads: list[dict[str, str | int]] = []
def fake_post(url: str, json: dict[str, str | int], timeout: int):
sent_payloads.append(json)
return httpx.Response(204, request=httpx.Request("POST", url))
monkeypatch.setattr(telemetry_service.httpx, "post", fake_post)
now = datetime(2026, 7, 13, 0, 0, 0)
assert CommunityTelemetryService.report_heartbeat(session=sqlite_session, now=now) is True
saved_setup = sqlite_session.scalar(select(DifySetup))
assert saved_setup is not None
assert saved_setup.install_reported_at is not None
assert saved_setup.last_heartbeat_at == now
assert [(payload["event"], payload["version"]) for payload in sent_payloads] == [
("install", "installed-version"),
("heartbeat", "running-version"),
]
@pytest.mark.parametrize("sqlite_session", [(DifySetup,)], indirect=True)
def test_report_heartbeat_skips_when_already_sent_today(
sqlite_session: Session, telemetry_enabled, monkeypatch: pytest.MonkeyPatch
):
setup = DifySetup(
version="1.0.0",
instance_id="d246c3a1-350b-406c-92c7-6043df680758",
install_reported_at=datetime(2026, 7, 13, 8, 0, 0),
last_heartbeat_at=datetime(2026, 7, 13, 9, 0, 0),
)
sqlite_session.add(setup)
sqlite_session.commit()
post_mock = Mock()
monkeypatch.setattr(telemetry_service.httpx, "post", post_mock)
assert (
CommunityTelemetryService.report_heartbeat(session=sqlite_session, now=datetime(2026, 7, 13, 12, 0, 0)) is False
)
post_mock.assert_not_called()
@pytest.mark.parametrize("sqlite_session", [(DifySetup,)], indirect=True)
def test_report_heartbeat_failure_does_not_mark_the_day_reported(
sqlite_session: Session, telemetry_enabled, monkeypatch: pytest.MonkeyPatch
):
setup = DifySetup(
version="1.0.0",
instance_id="d246c3a1-350b-406c-92c7-6043df680758",
install_reported_at=datetime(2026, 7, 13, 8, 0, 0),
)
sqlite_session.add(setup)
sqlite_session.commit()
def fake_post(url: str, json: dict[str, str | int], timeout: int):
raise httpx.ConnectError("offline", request=httpx.Request("POST", url))
monkeypatch.setattr(telemetry_service.httpx, "post", fake_post)
assert (
CommunityTelemetryService.report_heartbeat(session=sqlite_session, now=datetime(2026, 7, 13, 12, 0, 0)) is False
)
assert setup.last_heartbeat_at is None
def test_send_event_skips_when_telemetry_is_disabled(monkeypatch: pytest.MonkeyPatch):
monkeypatch.setattr(telemetry_service.dify_config, "DISABLE_TELEMETRY", True)
post_mock = Mock()
monkeypatch.setattr(telemetry_service.httpx, "post", post_mock)
assert CommunityTelemetryService._send_event({"event": "heartbeat"}) is False
post_mock.assert_not_called()
def test_send_event_skips_an_empty_fallback_endpoint(telemetry_enabled, monkeypatch: pytest.MonkeyPatch):
monkeypatch.setattr(telemetry_service.dify_config, "TELEMETRY_FALLBACK_ENDPOINT", "")
def fake_post(url: str, json: dict[str, str], timeout: int):
raise httpx.ConnectError("offline", request=httpx.Request("POST", url))
monkeypatch.setattr(telemetry_service.httpx, "post", fake_post)
assert CommunityTelemetryService._send_event({"event": "heartbeat"}) is False
def test_send_event_does_not_retry_the_same_endpoint(telemetry_enabled, monkeypatch: pytest.MonkeyPatch):
monkeypatch.setattr(
telemetry_service.dify_config,
"TELEMETRY_FALLBACK_ENDPOINT",
telemetry_service.dify_config.TELEMETRY_ENDPOINT,
)
post_mock = Mock(
return_value=httpx.Response(
204,
request=httpx.Request("POST", telemetry_service.dify_config.TELEMETRY_ENDPOINT),
)
)
monkeypatch.setattr(telemetry_service.httpx, "post", post_mock)
assert CommunityTelemetryService._send_event({"event": "heartbeat"}) is True
post_mock.assert_called_once()
def test_heartbeat_is_not_due_without_instance_id():
setup = DifySetup(version="1.0.0")
assert CommunityTelemetryService._is_heartbeat_due(setup, datetime(2026, 7, 13, 12, 0, 0)) is False
@pytest.mark.parametrize(
("value", "expected"),
[
("Linux", "linux"),
("Plan9", "unknown"),
],
)
def test_normalize_os(value: str, expected: str):
assert CommunityTelemetryService._normalize_os(value) == expected
@pytest.mark.parametrize(
("value", "expected"),
[
("x86_64", "amd64"),
("aarch64", "arm64"),
("armv7l", "arm"),
("i686", "386"),
("riscv64", "unknown"),
],
)
def test_normalize_arch(value: str, expected: str):
assert CommunityTelemetryService._normalize_arch(value) == expected
@@ -673,118 +673,3 @@ def test_dummy_variable_truncator_methods():
assert isinstance(result, TruncationResult)
assert result.result == segment
assert result.truncated is False
# ---------------------------------------------------------------------------
# Regression tests for langgenius/dify#39218.
#
# Before the fix, ``_truncate_array`` had a "Dirty fix" branch that
# unconditionally appended every ``File`` element to ``truncated_value``
# *before* the count cap and the byte-budget check, and *before*
# ``used_size`` was ever incremented. That made ``list[File]`` arrays:
# 1. uncapped by ``array_element_limit``,
# 2. uncounted against ``max_size_bytes``, and
# 3. always reported ``truncated=False``.
# The fix routes ``File`` through ``_truncate_json_primitives``'s dedicated
# ``File`` branch, which returns the file as-is with its real serialized
# size, while preserving the count cap and the byte budget.
# ---------------------------------------------------------------------------
class TestFileArrayTruncationRegression39218:
"""``list[File]`` must respect ``array_element_limit`` and the byte budget."""
@pytest.fixture
def truncator(self) -> VariableTruncator:
return VariableTruncator(
array_element_limit=3,
max_size_bytes=1000,
string_length_limit=50,
)
@staticmethod
def _make_file(name: str = "f") -> File:
return File(
id=name,
type=FileType.DOCUMENT,
transfer_method=FileTransferMethod.REMOTE_URL,
remote_url=f"https://example.com/{name}.txt",
filename=f"{name}.txt",
extension=".txt",
mime_type="text/plain",
size=1024,
)
def test_file_array_respects_element_count_cap(self, truncator: VariableTruncator) -> None:
# Use a target_size larger than ``count * file_size`` so the byte
# budget never binds — only the count cap should fire.
# Each File serializes to ~237 bytes; 3 files = ~713 bytes.
files = [self._make_file(f"f{i}") for i in range(500)]
result = truncator._truncate_array(files, target_size=10_000_000)
# Before the fix, all 500 File entries survived (``len(value)==500``,
# ``truncated==False``). After the fix, the array is capped at
# ``array_element_limit=3`` and ``truncated`` flips to True.
assert len(result.value) == 3
assert result.truncated is True
def test_file_array_reports_real_used_size(self, truncator: VariableTruncator) -> None:
# Large budget so the count cap fires before the byte budget does.
files = [self._make_file(f"f{i}") for i in range(500)]
result = truncator._truncate_array(files, target_size=10_000_000)
# Before the fix, ``used_size`` for a File array was the empty-array
# baseline of 2 bytes (``[]``), regardless of how many File entries
# actually returned. After the fix, ``used_size`` reflects the real
# serialized size of the returned ``File`` payload.
assert result.value_size > 100
assert result.truncated is True
def test_file_array_respects_byte_budget(self, truncator: VariableTruncator) -> None:
# Use a small ``target_size`` so the byte budget is the binding
# constraint. Each File serializes to ~237 bytes, so even one File
# blows the 200-byte budget.
files = [self._make_file(f"f{i}") for i in range(50)]
result = truncator._truncate_array(files, target_size=200)
# Before the fix, all 50 File entries survived and ``used_size``
# reported ``2`` (the empty-array baseline). After the fix, the
# loop sees the File payload: ``value_size`` reflects the real
# serialized size, and the loop stops after the first File because
# adding the next one would exceed ``target_size``.
assert len(result.value) == 1
assert result.value_size > 100 # the File's real serialized size
assert result.value_size <= 250 # in the ballpark of the budget
def test_mixed_array_counts_files_toward_cap(self, truncator: VariableTruncator) -> None:
mixed: list[object] = [
self._make_file("f0"),
"a",
self._make_file("f1"),
"b",
self._make_file("f2"),
"c",
self._make_file("f3"),
"d",
]
result = truncator._truncate_array(mixed, target_size=10_000_000)
# 8 items, cap of 3 → exactly 3 items. Files and primitives are
# counted together toward the cap.
assert len(result.value) == 3
assert result.truncated is True
def test_single_file_in_array_is_preserved(self, truncator: VariableTruncator) -> None:
result = truncator._truncate_array([self._make_file("only")], target_size=10_000_000)
# The File itself is not truncated — the dedicated ``File`` branch
# in ``_truncate_json_primitives`` returns the file untouched. Only
# the array-shape accounting changes.
assert len(result.value) == 1
assert isinstance(result.value[0], File)
assert result.value[0].id == "only"
assert result.truncated is False
File diff suppressed because it is too large Load Diff
@@ -1,16 +1,11 @@
"""Workflow-run service tests with real SQLite-bound session factories."""
from decimal import Decimal
from types import SimpleNamespace
from typing import Any, cast
from unittest.mock import MagicMock
import pytest
from sqlalchemy import Engine, event
from sqlalchemy.orm import Session, sessionmaker
from sqlalchemy import Engine
from models import Account, App, EndUser, Message, WorkflowRunTriggeredFrom
from models.enums import ConversationFromSource
from models import Account, App, EndUser, WorkflowRunTriggeredFrom
from services import workflow_run_service as service_module
from services.workflow_run_service import WorkflowRunService
@@ -27,11 +22,6 @@ def repository_factory_mocks(monkeypatch: pytest.MonkeyPatch) -> tuple[MagicMock
return node_repo, workflow_run_repo, factory
@pytest.fixture
def sqlalchemy_session_factory(sqlite_engine: Engine) -> sessionmaker[Session]:
return sessionmaker(bind=sqlite_engine, expire_on_commit=False)
def _app_model(**kwargs: Any) -> App:
return cast(App, SimpleNamespace(**kwargs))
@@ -44,22 +34,13 @@ def _end_user(**kwargs: Any) -> EndUser:
return cast(EndUser, SimpleNamespace(**kwargs))
def _message(*, message_id: str, workflow_run_id: str, conversation_id: str) -> Message:
message = Message(
app_id="app-1",
conversation_id=conversation_id,
query="query",
message={"role": "user", "content": "query"},
answer="answer",
message_unit_price=Decimal("0.0001"),
answer_unit_price=Decimal("0.0001"),
currency="USD",
from_source=ConversationFromSource.API,
)
message.id = message_id
message._inputs = {}
message.workflow_run_id = workflow_run_id
return message
def _fake_session_factory_returning_messages(messages: list[Any]) -> tuple[MagicMock, MagicMock]:
"""Build a session factory whose session returns the given messages."""
session = MagicMock()
session.scalars.return_value.all.return_value = messages
session_factory = MagicMock()
session_factory.return_value.__enter__.return_value = session
return session_factory, session
class TestWorkflowRunServiceInitialization:
@@ -67,51 +48,59 @@ class TestWorkflowRunServiceInitialization:
self,
monkeypatch: pytest.MonkeyPatch,
repository_factory_mocks: tuple[MagicMock, MagicMock, Any],
sqlite_engine: Engine,
) -> None:
monkeypatch.setattr(service_module, "db", SimpleNamespace(engine=sqlite_engine))
session_factory = MagicMock(name="session_factory")
sessionmaker_mock = MagicMock(return_value=session_factory)
monkeypatch.setattr(service_module, "sessionmaker", sessionmaker_mock)
monkeypatch.setattr(service_module, "db", SimpleNamespace(engine="db-engine"))
service = WorkflowRunService()
assert isinstance(service._session_factory, sessionmaker)
assert service._session_factory.kw["bind"] is sqlite_engine
assert service._session_factory.kw["expire_on_commit"] is False
sessionmaker_mock.assert_called_once_with(bind="db-engine", expire_on_commit=False)
assert service._session_factory is session_factory
def test___init___should_create_sessionmaker_when_engine_is_provided(
self,
monkeypatch: pytest.MonkeyPatch,
repository_factory_mocks: tuple[MagicMock, MagicMock, Any],
sqlite_engine: Engine,
) -> None:
service = WorkflowRunService(session_factory=sqlite_engine)
class FakeEngine:
pass
assert isinstance(service._session_factory, sessionmaker)
assert service._session_factory.kw["bind"] is sqlite_engine
assert service._session_factory.kw["expire_on_commit"] is False
session_factory = MagicMock(name="session_factory")
sessionmaker_mock = MagicMock(return_value=session_factory)
monkeypatch.setattr(service_module, "Engine", FakeEngine)
monkeypatch.setattr(service_module, "sessionmaker", sessionmaker_mock)
engine = cast(Engine, FakeEngine())
service = WorkflowRunService(session_factory=engine)
sessionmaker_mock.assert_called_once_with(bind=engine, expire_on_commit=False)
assert service._session_factory is session_factory
def test___init___should_keep_provided_sessionmaker_and_create_repositories(
self,
repository_factory_mocks: tuple[MagicMock, MagicMock, Any],
sqlalchemy_session_factory: sessionmaker[Session],
) -> None:
node_repo, workflow_run_repo, factory = repository_factory_mocks
session_factory = MagicMock(name="session_factory")
service = WorkflowRunService(session_factory=sqlalchemy_session_factory)
service = WorkflowRunService(session_factory=session_factory)
assert service._session_factory is sqlalchemy_session_factory
assert service._session_factory is session_factory
assert service._node_execution_service_repo is node_repo
assert service._workflow_run_repo is workflow_run_repo
factory.create_api_workflow_node_execution_repository.assert_called_once_with(sqlalchemy_session_factory)
factory.create_api_workflow_run_repository.assert_called_once_with(sqlalchemy_session_factory)
factory.create_api_workflow_node_execution_repository.assert_called_once_with(session_factory)
factory.create_api_workflow_run_repository.assert_called_once_with(session_factory)
class TestWorkflowRunServiceQueries:
def test_get_paginate_workflow_runs_should_forward_filters_and_parse_limit(
self,
repository_factory_mocks: tuple[MagicMock, MagicMock, Any],
sqlalchemy_session_factory: sessionmaker[Session],
) -> None:
_, workflow_run_repo, _ = repository_factory_mocks
service = WorkflowRunService(session_factory=sqlalchemy_session_factory)
service = WorkflowRunService(session_factory=MagicMock(name="session_factory"))
app_model = _app_model(tenant_id="tenant-1", id="app-1")
expected = MagicMock(name="pagination")
workflow_run_repo.get_paginated_workflow_runs.return_value = expected
@@ -133,24 +122,20 @@ class TestWorkflowRunServiceQueries:
status="succeeded",
)
@pytest.mark.parametrize("sqlite_session", [(Message,)], indirect=True)
def test_get_paginate_advanced_chat_workflow_runs_should_attach_message_fields_when_message_exists(
self,
repository_factory_mocks: tuple[MagicMock, MagicMock, Any],
monkeypatch: pytest.MonkeyPatch,
sqlalchemy_session_factory: sessionmaker[Session],
sqlite_session: Session,
) -> None:
service = WorkflowRunService(session_factory=sqlalchemy_session_factory)
message = SimpleNamespace(id="msg-1", conversation_id="conv-1", workflow_run_id="run-1")
session_factory, session = _fake_session_factory_returning_messages([message])
service = WorkflowRunService(session_factory=session_factory)
app_model = _app_model(tenant_id="tenant-1", id="app-1")
run_with_message = SimpleNamespace(id="run-1", status="running")
run_without_message = SimpleNamespace(id="run-2", status="succeeded")
pagination = SimpleNamespace(data=[run_with_message, run_without_message])
monkeypatch.setattr(service, "get_paginate_workflow_runs", MagicMock(return_value=pagination))
sqlite_session.add(_message(message_id="msg-1", conversation_id="conv-1", workflow_run_id="run-1"))
sqlite_session.commit()
result = service.get_paginate_advanced_chat_workflow_runs(app_model=app_model, args={"limit": "2"})
assert result is pagination
@@ -160,49 +145,39 @@ class TestWorkflowRunServiceQueries:
assert result.data[0].status == "running"
assert not hasattr(result.data[1], "message_id")
assert result.data[1].id == "run-2"
# Messages are batch-loaded in a single query, not one per run.
session_factory.assert_called_once_with()
session.scalars.assert_called_once()
@pytest.mark.parametrize("sqlite_session", [(Message,)], indirect=True)
def test_get_paginate_advanced_chat_workflow_runs_batch_loads_messages_without_n_plus_one(
self,
repository_factory_mocks: tuple[MagicMock, MagicMock, Any],
monkeypatch: pytest.MonkeyPatch,
sqlalchemy_session_factory: sessionmaker[Session],
sqlite_session: Session,
) -> None:
"""Messages must load with a constant query count regardless of run count.
Previously the deprecated WorkflowRun.message property issued one query per
run (N+1); they are now batch-loaded in a single query.
"""
service = WorkflowRunService(session_factory=sqlalchemy_session_factory)
session_factory, session = _fake_session_factory_returning_messages([])
service = WorkflowRunService(session_factory=session_factory)
app_model = _app_model(tenant_id="tenant-1", id="app-1")
runs = [SimpleNamespace(id=f"run-{i}", status="succeeded") for i in range(5)]
pagination = SimpleNamespace(data=runs)
monkeypatch.setattr(service, "get_paginate_workflow_runs", MagicMock(return_value=pagination))
message_query_count = 0
service.get_paginate_advanced_chat_workflow_runs(app_model=app_model, args={})
def count_message_query(*_args: object) -> None:
nonlocal message_query_count
message_query_count += 1
engine = sqlite_session.get_bind()
event.listen(engine, "before_cursor_execute", count_message_query)
try:
service.get_paginate_advanced_chat_workflow_runs(app_model=app_model, args={})
finally:
event.remove(engine, "before_cursor_execute", count_message_query)
assert all(not hasattr(run, "message_id") for run in runs)
assert message_query_count == 1
# Exactly one message query for the whole page, independent of run count.
session_factory.assert_called_once_with()
assert session.scalars.call_count == 1
def test_get_workflow_run_should_delegate_to_repository_by_tenant_and_app(
self,
repository_factory_mocks: tuple[MagicMock, MagicMock, Any],
sqlalchemy_session_factory: sessionmaker[Session],
) -> None:
_, workflow_run_repo, _ = repository_factory_mocks
service = WorkflowRunService(session_factory=sqlalchemy_session_factory)
service = WorkflowRunService(session_factory=MagicMock(name="session_factory"))
app_model = _app_model(tenant_id="tenant-1", id="app-1")
expected = MagicMock(name="workflow_run")
workflow_run_repo.get_workflow_run_by_id.return_value = expected
@@ -219,10 +194,9 @@ class TestWorkflowRunServiceQueries:
def test_get_workflow_runs_count_should_forward_optional_filters(
self,
repository_factory_mocks: tuple[MagicMock, MagicMock, Any],
sqlalchemy_session_factory: sessionmaker[Session],
) -> None:
_, workflow_run_repo, _ = repository_factory_mocks
service = WorkflowRunService(session_factory=sqlalchemy_session_factory)
service = WorkflowRunService(session_factory=MagicMock(name="session_factory"))
app_model = _app_model(tenant_id="tenant-1", id="app-1")
expected = {"total": 3, "succeeded": 2}
workflow_run_repo.get_workflow_runs_count.return_value = expected
@@ -247,9 +221,8 @@ class TestWorkflowRunServiceQueries:
self,
repository_factory_mocks: tuple[MagicMock, MagicMock, Any],
monkeypatch: pytest.MonkeyPatch,
sqlalchemy_session_factory: sessionmaker[Session],
) -> None:
service = WorkflowRunService(session_factory=sqlalchemy_session_factory)
service = WorkflowRunService(session_factory=MagicMock(name="session_factory"))
monkeypatch.setattr(service, "get_workflow_run", MagicMock(return_value=None))
app_model = _app_model(id="app-1")
user = _account(current_tenant_id="tenant-1")
@@ -262,10 +235,9 @@ class TestWorkflowRunServiceQueries:
self,
repository_factory_mocks: tuple[MagicMock, MagicMock, Any],
monkeypatch: pytest.MonkeyPatch,
sqlalchemy_session_factory: sessionmaker[Session],
) -> None:
node_repo, _, _ = repository_factory_mocks
service = WorkflowRunService(session_factory=sqlalchemy_session_factory)
service = WorkflowRunService(session_factory=MagicMock(name="session_factory"))
monkeypatch.setattr(service, "get_workflow_run", MagicMock(return_value=SimpleNamespace(id="run-1")))
class FakeEndUser:
@@ -295,10 +267,9 @@ class TestWorkflowRunServiceQueries:
self,
repository_factory_mocks: tuple[MagicMock, MagicMock, Any],
monkeypatch: pytest.MonkeyPatch,
sqlalchemy_session_factory: sessionmaker[Session],
) -> None:
node_repo, _, _ = repository_factory_mocks
service = WorkflowRunService(session_factory=sqlalchemy_session_factory)
service = WorkflowRunService(session_factory=MagicMock(name="session_factory"))
monkeypatch.setattr(service, "get_workflow_run", MagicMock(return_value=SimpleNamespace(id="run-1")))
app_model = _app_model(id="app-1")
user = _account(current_tenant_id="tenant-account")
@@ -322,9 +293,8 @@ class TestWorkflowRunServiceQueries:
self,
repository_factory_mocks: tuple[MagicMock, MagicMock, Any],
monkeypatch: pytest.MonkeyPatch,
sqlalchemy_session_factory: sessionmaker[Session],
) -> None:
service = WorkflowRunService(session_factory=sqlalchemy_session_factory)
service = WorkflowRunService(session_factory=MagicMock(name="session_factory"))
monkeypatch.setattr(service, "get_workflow_run", MagicMock(return_value=SimpleNamespace(id="run-1")))
app_model = _app_model(id="app-1")
user = _account(current_tenant_id=None)
@@ -1,73 +1,32 @@
"""Simplified unit tests for DraftVarLoader focusing on core functionality."""
import json
from datetime import datetime
from unittest.mock import Mock, patch
import pytest
from sqlalchemy import Engine
from sqlalchemy.orm import Session
from core.workflow.file_reference import build_file_reference
from extensions.storage.storage_type import StorageType
from graphon.file import File, FileTransferMethod, FileType
from graphon.variables.segments import ObjectSegment, StringSegment
from graphon.variables.types import SegmentType
from models.enums import CreatorUserRole
from models.model import UploadFile
from models.workflow import WorkflowDraftVariable, WorkflowDraftVariableFile
from services.workflow_draft_variable_service import DraftVarLoader
def _persist_offloaded_variable(
sqlite_session: Session,
*,
node_id: str,
name: str,
) -> WorkflowDraftVariable:
upload_file = UploadFile(
tenant_id="test-tenant-id",
storage_type=StorageType.LOCAL,
key=f"storage/key/{name}.txt",
name=f"{name}.txt",
size=10,
extension=".txt",
mime_type="text/plain",
created_by_role=CreatorUserRole.ACCOUNT,
created_by="test-user-id",
created_at=datetime(2025, 1, 1),
used=True,
)
variable_file = WorkflowDraftVariableFile(
tenant_id="test-tenant-id",
app_id="test-app-id",
user_id="test-user-id",
upload_file_id=upload_file.id,
size=10,
length=None,
value_type=SegmentType.STRING,
)
draft_variable = WorkflowDraftVariable.new_node_variable(
app_id="test-app-id",
user_id="test-user-id",
node_id=node_id,
name=name,
value=StringSegment(value="truncated"),
node_execution_id=f"execution-{node_id}",
file_id=variable_file.id,
)
sqlite_session.add_all([upload_file, variable_file, draft_variable])
return draft_variable
class TestDraftVarLoaderSimple:
"""Simplified unit tests for DraftVarLoader core methods."""
@pytest.fixture
def draft_var_loader(self, sqlite_engine: Engine):
def mock_engine(self) -> Engine:
return Mock(spec=Engine)
@pytest.fixture
def draft_var_loader(self, mock_engine):
"""Create DraftVarLoader instance for testing."""
return DraftVarLoader(
engine=sqlite_engine,
engine=mock_engine,
app_id="test-app-id",
tenant_id="test-tenant-id",
user_id="test-user-id",
@@ -246,109 +205,131 @@ class TestDraftVarLoaderSimple:
assert variable.value == rebuilt_file
rebuild_file.assert_called_once_with(file_mapping=raw_file, tenant_id="tenant-1")
@pytest.mark.parametrize(
"sqlite_session",
[(WorkflowDraftVariable, WorkflowDraftVariableFile, UploadFile)],
indirect=True,
)
def test_load_variables_with_offloaded_variables_unit(
self,
draft_var_loader: DraftVarLoader,
sqlite_session: Session,
):
def test_load_variables_with_offloaded_variables_unit(self, draft_var_loader):
"""Test load_variables method with mix of regular and offloaded variables."""
selectors = [["node1", "regular_var"], ["node2", "offloaded_var"]]
regular_draft_var = WorkflowDraftVariable.new_node_variable(
app_id="test-app-id",
user_id="test-user-id",
node_id="node1",
name="regular_var",
value=StringSegment(value="regular_value"),
node_execution_id="execution-node1",
)
# Mock regular variable
regular_draft_var = Mock(spec=WorkflowDraftVariable)
regular_draft_var.is_truncated.return_value = False
regular_draft_var.node_id = "node1"
regular_draft_var.name = "regular_var"
regular_draft_var.get_value.return_value = StringSegment(value="regular_value")
regular_draft_var.get_selector.return_value = ["node1", "regular_var"]
regular_draft_var.id = "regular-var-id"
regular_draft_var.description = "regular description"
offloaded_draft_var = _persist_offloaded_variable(
sqlite_session,
node_id="node2",
name="offloaded_var",
)
distractor = WorkflowDraftVariable.new_node_variable(
app_id="test-app-id",
user_id="another-user",
node_id="node1",
name="regular_var",
value=StringSegment(value="wrong user"),
node_execution_id="execution-distractor",
)
sqlite_session.add_all([regular_draft_var, distractor])
sqlite_session.commit()
offloaded_variable = Mock()
offloaded_variable.id = offloaded_draft_var.id
offloaded_variable.selector = ["node2", "offloaded_var"]
# Mock offloaded variable
upload_file = Mock(spec=UploadFile)
upload_file.key = "storage/key/offloaded.txt"
with (
patch("services.workflow_draft_variable_service.StorageKeyLoader"),
patch.object(
draft_var_loader,
"_load_offloaded_variable",
return_value=(("node2", "offloaded_var"), offloaded_variable),
) as load_offloaded,
patch("services.workflow_draft_variable_service.ThreadPoolExecutor") as executor_cls,
):
executor = executor_cls.return_value.__enter__.return_value
executor.map.side_effect = lambda function, values: [function(value) for value in values]
variable_file = Mock(spec=WorkflowDraftVariableFile)
variable_file.value_type = SegmentType.STRING
variable_file.upload_file = upload_file
result = draft_var_loader.load_variables(selectors)
offloaded_draft_var = Mock(spec=WorkflowDraftVariable)
offloaded_draft_var.is_truncated.return_value = True
offloaded_draft_var.node_id = "node2"
offloaded_draft_var.name = "offloaded_var"
offloaded_draft_var.get_selector.return_value = ["node2", "offloaded_var"]
offloaded_draft_var.variable_file = variable_file
offloaded_draft_var.id = "offloaded-var-id"
offloaded_draft_var.description = "offloaded description"
assert {variable.id for variable in result} == {regular_draft_var.id, offloaded_draft_var.id}
load_offloaded.assert_called_once()
loaded_offloaded = load_offloaded.call_args.args[0]
assert isinstance(loaded_offloaded, WorkflowDraftVariable)
assert loaded_offloaded.id == offloaded_draft_var.id
assert loaded_offloaded.variable_file is not None
assert loaded_offloaded.variable_file.upload_file is not None
assert loaded_offloaded.variable_file.upload_file.key == "storage/key/offloaded_var.txt"
draft_vars = [regular_draft_var, offloaded_draft_var]
@pytest.mark.parametrize(
"sqlite_session",
[(WorkflowDraftVariable, WorkflowDraftVariableFile, UploadFile)],
indirect=True,
)
def test_load_variables_all_offloaded_variables_unit(
self,
draft_var_loader: DraftVarLoader,
sqlite_session: Session,
):
with patch("services.workflow_draft_variable_service.Session") as mock_session_cls:
mock_session = Mock()
mock_session_cls.return_value.__enter__.return_value = mock_session
mock_service = Mock()
mock_service.get_draft_variables_by_selectors.return_value = draft_vars
with patch(
"services.workflow_draft_variable_service.WorkflowDraftVariableService", return_value=mock_service
):
with patch("services.workflow_draft_variable_service.StorageKeyLoader"):
with patch("factories.variable_factory.segment_to_variable") as mock_segment_to_variable:
# Mock regular variable creation
regular_variable = Mock()
regular_variable.selector = ["node1", "regular_var"]
# Mock offloaded variable creation
offloaded_variable = Mock()
offloaded_variable.selector = ["node2", "offloaded_var"]
mock_segment_to_variable.return_value = regular_variable
with patch("services.workflow_draft_variable_service.storage") as mock_storage:
mock_storage.load.return_value = b"offloaded_content"
with patch.object(draft_var_loader, "_load_offloaded_variable") as mock_load_offloaded:
mock_load_offloaded.return_value = (("node2", "offloaded_var"), offloaded_variable)
with patch("concurrent.futures.ThreadPoolExecutor") as mock_executor_cls:
mock_executor = Mock()
mock_executor_cls.return_value.__enter__.return_value = mock_executor
mock_executor.map.return_value = [(("node2", "offloaded_var"), offloaded_variable)]
# Execute the method
result = draft_var_loader.load_variables(selectors)
# Verify results
assert len(result) == 2
# Verify service method was called
mock_service.get_draft_variables_by_selectors.assert_called_once_with(
draft_var_loader._app_id,
selectors,
user_id=draft_var_loader._user_id,
)
# Verify offloaded variable loading was called
mock_load_offloaded.assert_called_once_with(offloaded_draft_var)
def test_load_variables_all_offloaded_variables_unit(self, draft_var_loader):
"""Test load_variables method with only offloaded variables."""
selectors = [["node1", "offloaded_var1"], ["node2", "offloaded_var2"]]
offloaded_var1 = _persist_offloaded_variable(
sqlite_session,
node_id="node1",
name="offloaded_var1",
)
offloaded_var2 = _persist_offloaded_variable(
sqlite_session,
node_id="node2",
name="offloaded_var2",
)
sqlite_session.commit()
with (
patch("services.workflow_draft_variable_service.StorageKeyLoader"),
patch("services.workflow_draft_variable_service.ThreadPoolExecutor") as executor_cls,
):
executor = executor_cls.return_value.__enter__.return_value
executor.map.return_value = [
(("node1", "offloaded_var1"), Mock()),
(("node2", "offloaded_var2"), Mock()),
]
# Mock first offloaded variable
offloaded_var1 = Mock(spec=WorkflowDraftVariable)
offloaded_var1.is_truncated.return_value = True
offloaded_var1.node_id = "node1"
offloaded_var1.name = "offloaded_var1"
result = draft_var_loader.load_variables(selectors)
# Mock second offloaded variable
offloaded_var2 = Mock(spec=WorkflowDraftVariable)
offloaded_var2.is_truncated.return_value = True
offloaded_var2.node_id = "node2"
offloaded_var2.name = "offloaded_var2"
assert len(result) == 2
executor_cls.assert_called_once_with(max_workers=10)
executor.map.assert_called_once()
loaded_draft_vars = executor.map.call_args.args[1]
assert {variable.id for variable in loaded_draft_vars} == {offloaded_var1.id, offloaded_var2.id}
assert all(variable.variable_file.upload_file is not None for variable in loaded_draft_vars)
draft_vars = [offloaded_var1, offloaded_var2]
with patch("services.workflow_draft_variable_service.Session") as mock_session_cls:
mock_session = Mock()
mock_session_cls.return_value.__enter__.return_value = mock_session
mock_service = Mock()
mock_service.get_draft_variables_by_selectors.return_value = draft_vars
with patch(
"services.workflow_draft_variable_service.WorkflowDraftVariableService", return_value=mock_service
):
with patch("services.workflow_draft_variable_service.StorageKeyLoader"):
with patch("services.workflow_draft_variable_service.ThreadPoolExecutor") as mock_executor_cls:
mock_executor = Mock()
mock_executor_cls.return_value.__enter__.return_value = mock_executor
mock_executor.map.return_value = [
(("node1", "offloaded_var1"), Mock()),
(("node2", "offloaded_var2"), Mock()),
]
# Execute the method
result = draft_var_loader.load_variables(selectors)
# Verify results - since we have only offloaded variables, should have 2 results
assert len(result) == 2
# Verify ThreadPoolExecutor was used
mock_executor_cls.assert_called_once_with(max_workers=10)
mock_executor.map.assert_called_once()
@@ -1,40 +0,0 @@
from types import SimpleNamespace
from unittest.mock import MagicMock, Mock
import pytest
from tasks import community_telemetry_task
def _configure_task_session(monkeypatch: pytest.MonkeyPatch) -> Mock:
session = Mock()
session_factory = MagicMock()
session_factory.return_value.__enter__.return_value = session
monkeypatch.setattr(community_telemetry_task, "db", SimpleNamespace(engine=object()))
monkeypatch.setattr(community_telemetry_task, "sessionmaker", Mock(return_value=session_factory))
return session
def test_send_community_telemetry_heartbeat_reports_with_a_database_session(monkeypatch: pytest.MonkeyPatch):
session = _configure_task_session(monkeypatch)
report_heartbeat = Mock()
monkeypatch.setattr(community_telemetry_task.CommunityTelemetryService, "report_heartbeat", report_heartbeat)
community_telemetry_task.send_community_telemetry_heartbeat.run()
report_heartbeat.assert_called_once_with(session=session)
def test_send_community_telemetry_heartbeat_swallows_report_errors(monkeypatch: pytest.MonkeyPatch):
_configure_task_session(monkeypatch)
monkeypatch.setattr(
community_telemetry_task.CommunityTelemetryService,
"report_heartbeat",
Mock(side_effect=RuntimeError("telemetry unavailable")),
)
log_debug = Mock()
monkeypatch.setattr(community_telemetry_task.logger, "debug", log_debug)
community_telemetry_task.send_community_telemetry_heartbeat.run()
log_debug.assert_called_once_with("Failed to process community telemetry heartbeat", exc_info=True)
-314
View File
@@ -1,314 +0,0 @@
"""Enterprise license gating performed by the global ``before_request`` hook."""
from unittest.mock import patch
import pytest
from flask import Blueprint, Flask
from flask_restx import Resource
from app_factory import create_flask_app_with_configs
from libs.external_api import ExternalApi
from services.feature_service import LicenseStatus
INVALID_STATUSES = [LicenseStatus.INACTIVE, LicenseStatus.EXPIRED, LicenseStatus.LOST]
VALID_STATUSES = [LicenseStatus.ACTIVE, LicenseStatus.EXPIRING]
def _license(status: LicenseStatus | None):
return patch("app_factory.EnterpriseService.get_cached_license_status", return_value=status)
def _enterprise(enabled: bool = True):
return patch("app_factory.dify_config.ENTERPRISE_ENABLED", enabled)
@pytest.fixture
def gated_app() -> Flask:
app = create_flask_app_with_configs()
@app.route("/v1/chat-messages", methods=["POST"])
def service_api_route():
return {"surface": "service_api"}
@app.route("/v1/")
def service_api_index_route():
return {"surface": "service_api_index"}
@app.route("/mcp/server/<server_code>/mcp", methods=["POST"])
def mcp_route(server_code: str):
return {"surface": "mcp"}
@app.route("/triggers/webhook/<webhook_id>", methods=["POST"])
def trigger_route(webhook_id: str):
return {"surface": "triggers"}
@app.route("/console/api/apps")
def console_route():
return {"surface": "console"}
@app.route("/console/api/login", methods=["POST"])
def console_bootstrap_route():
return {"surface": "console_bootstrap"}
@app.route("/api/messages")
def webapp_route():
return {"surface": "webapp"}
@app.route("/api/system-features")
def webapp_bootstrap_route():
return {"surface": "webapp_bootstrap"}
@app.route("/health")
def health_route():
return {"surface": "health"}
@app.route("/inner/api/rbac/check-access", methods=["POST"])
def inner_api_route():
return {"surface": "inner_api"}
@app.route("/files/upload/for-plugin", methods=["POST"])
def files_route():
return {"surface": "files"}
return app
class TestServiceApiLicenseGate:
"""/v1 is a bearer-token surface, so it is gated with an opaque 403."""
@pytest.mark.parametrize("status", INVALID_STATUSES)
def test_blocks_when_license_invalid(self, gated_app: Flask, status: LicenseStatus):
with _enterprise(), _license(status):
response = gated_app.test_client().post("/v1/chat-messages")
assert response.status_code == 403
def test_block_response_carries_machine_readable_marker(self, gated_app: Flask):
with _enterprise(), _license(LicenseStatus.EXPIRED):
response = gated_app.test_client().post("/v1/chat-messages")
assert b"license_required" in response.data
def test_block_response_does_not_leak_license_status(self, gated_app: Flask):
with _enterprise(), _license(LicenseStatus.EXPIRED):
response = gated_app.test_client().post("/v1/chat-messages")
assert b"expired" not in response.data.lower()
def test_blocks_when_license_status_unavailable(self, gated_app: Flask):
with _enterprise(), _license(None):
response = gated_app.test_client().post("/v1/chat-messages")
assert response.status_code == 403
def test_blocks_when_license_lookup_raises(self, gated_app: Flask):
lookup_failed = patch(
"app_factory.EnterpriseService.get_cached_license_status",
side_effect=RuntimeError("enterprise api unreachable"),
)
with _enterprise(), lookup_failed:
response = gated_app.test_client().post("/v1/chat-messages")
assert response.status_code == 403
def test_blocks_index_route(self, gated_app: Flask):
"""/v1 has no sign-in page to bootstrap, so nothing on it is exempt."""
with _enterprise(), _license(LicenseStatus.EXPIRED):
response = gated_app.test_client().get("/v1/")
assert response.status_code == 403
@pytest.mark.parametrize("status", VALID_STATUSES)
def test_allows_when_license_valid(self, gated_app: Flask, status: LicenseStatus):
with _enterprise(), _license(status):
response = gated_app.test_client().post("/v1/chat-messages")
assert response.status_code == 200
def test_allows_unclassified_status(self, gated_app: Flask):
"""LicenseStatus.NONE is not in the blocked set — parity with console/webapp."""
with _enterprise(), _license(LicenseStatus.NONE):
response = gated_app.test_client().post("/v1/chat-messages")
assert response.status_code == 200
@pytest.mark.parametrize("status", INVALID_STATUSES)
def test_does_not_gate_community_edition(self, gated_app: Flask, status: LicenseStatus):
with _enterprise(False), _license(status):
response = gated_app.test_client().post("/v1/chat-messages")
assert response.status_code == 200
class TestMcpLicenseGate:
"""/mcp invokes apps for external MCP clients, so it is gated like the Service API."""
@pytest.mark.parametrize("status", INVALID_STATUSES)
def test_blocks_when_license_invalid(self, gated_app: Flask, status: LicenseStatus):
with _enterprise(), _license(status):
response = gated_app.test_client().post("/mcp/server/srv-code/mcp")
assert response.status_code == 403
def test_block_response_carries_machine_readable_marker(self, gated_app: Flask):
with _enterprise(), _license(LicenseStatus.EXPIRED):
response = gated_app.test_client().post("/mcp/server/srv-code/mcp")
assert b"license_required" in response.data
@pytest.mark.parametrize("status", VALID_STATUSES)
def test_allows_when_license_valid(self, gated_app: Flask, status: LicenseStatus):
with _enterprise(), _license(status):
response = gated_app.test_client().post("/mcp/server/srv-code/mcp")
assert response.status_code == 200
@pytest.mark.parametrize("status", INVALID_STATUSES)
def test_does_not_gate_community_edition(self, gated_app: Flask, status: LicenseStatus):
with _enterprise(False), _license(status):
response = gated_app.test_client().post("/mcp/server/srv-code/mcp")
assert response.status_code == 200
class TestTriggerLicenseGate:
"""Inbound webhooks are refused so senders retry, rather than dropping events."""
@pytest.mark.parametrize("status", INVALID_STATUSES)
def test_blocks_when_license_invalid(self, gated_app: Flask, status: LicenseStatus):
with _enterprise(), _license(status):
response = gated_app.test_client().post("/triggers/webhook/hook-id")
assert response.status_code == 503
def test_block_response_carries_machine_readable_marker(self, gated_app: Flask):
with _enterprise(), _license(LicenseStatus.EXPIRED):
response = gated_app.test_client().post("/triggers/webhook/hook-id")
assert b"license_required" in response.data
@pytest.mark.parametrize("status", VALID_STATUSES)
def test_allows_when_license_valid(self, gated_app: Flask, status: LicenseStatus):
with _enterprise(), _license(status):
response = gated_app.test_client().post("/triggers/webhook/hook-id")
assert response.status_code == 200
@pytest.mark.parametrize("status", INVALID_STATUSES)
def test_does_not_gate_community_edition(self, gated_app: Flask, status: LicenseStatus):
with _enterprise(False), _license(status):
response = gated_app.test_client().post("/triggers/webhook/hook-id")
assert response.status_code == 200
class TestGateThroughRealErrorHandlers:
"""Gate errors must survive each blueprint's error handling: flask-restx vs plain Flask."""
@pytest.fixture
def wired_app(self) -> Flask:
app = create_flask_app_with_configs()
service_api_bp = Blueprint("service_api_test", __name__, url_prefix="/v1")
api = ExternalApi(service_api_bp)
@api.route("/chat-messages")
class ChatMessages(Resource):
def post(self):
return {"surface": "service_api"}
app.register_blueprint(service_api_bp)
trigger_bp = Blueprint("trigger_test", __name__, url_prefix="/triggers")
@trigger_bp.route("/webhook/<webhook_id>", methods=["POST"])
def webhook_route(webhook_id: str):
return {"surface": "triggers"}
app.register_blueprint(trigger_bp)
return app
def test_service_api_block_is_json_with_license_marker(self, wired_app: Flask):
with _enterprise(), _license(LicenseStatus.EXPIRED):
response = wired_app.test_client().post("/v1/chat-messages")
assert response.status_code == 403
body = response.get_json()
assert body["message"] == "license_required"
assert body["status"] == 403
def test_service_api_block_does_not_clear_cookies(self, wired_app: Flask):
"""Force-logout cookie clearing belongs to the cookie-authed surfaces only."""
with _enterprise(), _license(LicenseStatus.EXPIRED):
response = wired_app.test_client().post("/v1/chat-messages")
assert response.headers.getlist("Set-Cookie") == []
def test_trigger_block_survives_plain_blueprint_handling(self, wired_app: Flask):
with _enterprise(), _license(LicenseStatus.EXPIRED):
response = wired_app.test_client().post("/triggers/webhook/hook-id")
assert response.status_code == 503
assert b"license_required" in response.data
def test_surfaces_are_reachable_when_license_valid(self, wired_app: Flask):
with _enterprise(), _license(LicenseStatus.ACTIVE):
service_api = wired_app.test_client().post("/v1/chat-messages")
triggers = wired_app.test_client().post("/triggers/webhook/hook-id")
assert service_api.status_code == 200
assert triggers.status_code == 200
class TestUngatedSurfaces:
"""Surfaces that must stay reachable while the license is invalid."""
def test_inner_api_is_not_gated(self, gated_app: Flask):
"""dify-enterprise control plane — gating it could block license recovery itself."""
with _enterprise(), _license(LicenseStatus.EXPIRED):
response = gated_app.test_client().post("/inner/api/rbac/check-access")
assert response.status_code == 200
def test_files_data_plane_is_not_gated(self, gated_app: Flask):
"""Signed file URLs are fetched by the plugin daemon and by LLM vendors."""
with _enterprise(), _license(LicenseStatus.EXPIRED):
response = gated_app.test_client().post("/files/upload/for-plugin")
assert response.status_code == 200
class TestSessionSurfaceLicenseGate:
"""Console and webapp are cookie-authed, so they keep force-logout 401 semantics."""
@pytest.mark.parametrize("status", INVALID_STATUSES)
def test_blocks_console_with_force_logout(self, gated_app: Flask, status: LicenseStatus):
with _enterprise(), _license(status):
response = gated_app.test_client().get("/console/api/apps")
assert response.status_code == 401
@pytest.mark.parametrize("status", INVALID_STATUSES)
def test_blocks_webapp_with_force_logout(self, gated_app: Flask, status: LicenseStatus):
with _enterprise(), _license(status):
response = gated_app.test_client().get("/api/messages")
assert response.status_code == 401
def test_console_bootstrap_route_stays_reachable(self, gated_app: Flask):
with _enterprise(), _license(LicenseStatus.EXPIRED):
response = gated_app.test_client().post("/console/api/login")
assert response.status_code == 200
def test_webapp_bootstrap_route_stays_reachable(self, gated_app: Flask):
with _enterprise(), _license(LicenseStatus.EXPIRED):
response = gated_app.test_client().get("/api/system-features")
assert response.status_code == 200
def test_health_route_is_never_gated(self, gated_app: Flask):
with _enterprise(), _license(LicenseStatus.EXPIRED):
response = gated_app.test_client().get("/health")
assert response.status_code == 200
@@ -1,83 +0,0 @@
from concurrent.futures import ThreadPoolExecutor
from pathlib import Path
from threading import Barrier
import pytest
from sqlalchemy import create_engine, inspect, text
from sqlalchemy.engine import URL, Engine
from sqlalchemy.orm import Session, sessionmaker
from sqlalchemy.pool import QueuePool
import core.db.session_factory as session_factory_module
from models.account import Account
from models.base import TypeBase
from models.model import ExporleBanner
def test_sqlite_session_contains_the_full_registered_schema(sqlite_session: Session) -> None:
table_names = set(inspect(sqlite_session.get_bind()).get_table_names())
assert table_names == set(TypeBase.metadata.tables)
@pytest.mark.parametrize("sqlite_session", [(Account,)], indirect=True)
def test_sqlite_session_accepts_deferred_legacy_indirect_parameters(sqlite_session: Session) -> None:
"""Prove legacy model parameters no longer limit the copied schema."""
assert inspect(sqlite_session.get_bind()).has_table(ExporleBanner.__tablename__)
def test_sqlite_engine_is_a_pristine_file_copy(
sqlite_engine: Engine,
request: pytest.FixtureRequest,
) -> None:
sqlite_database_template: Path = request.getfixturevalue("_sqlite_database_template")
assert isinstance(sqlite_engine.pool, QueuePool)
assert sqlite_engine.url.database != str(sqlite_database_template)
with sqlite_engine.begin() as connection:
connection.execute(text("CREATE TABLE per_test_mutation (value INTEGER NOT NULL)"))
template_engine = create_engine(URL.create("sqlite", database=str(sqlite_database_template)))
try:
assert not inspect(template_engine).has_table("per_test_mutation")
finally:
template_engine.dispose()
def test_core_session_factory_uses_the_shared_sqlite_session_factory(
sqlite_session_factory: sessionmaker[Session],
) -> None:
assert session_factory_module.session_factory.get_session_maker() is sqlite_session_factory
with sqlite_session_factory.begin() as session:
session.execute(text("CREATE TABLE global_factory_probe (value INTEGER NOT NULL)"))
session.execute(text("INSERT INTO global_factory_probe (value) VALUES (42)"))
with session_factory_module.session_factory.create_session() as session:
assert session.scalar(text("SELECT value FROM global_factory_probe")) == 42
def test_sqlite_session_factory_shares_one_database_across_worker_sessions(
sqlite_session_factory: sessionmaker[Session],
) -> None:
with sqlite_session_factory.begin() as session:
session.execute(text("CREATE TABLE thread_probe (value INTEGER NOT NULL)"))
session.execute(text("INSERT INTO thread_probe (value) VALUES (42)"))
worker_barrier = Barrier(2)
def read_value() -> tuple[int, int]:
with sqlite_session_factory() as session:
connection = session.connection()
worker_barrier.wait(timeout=1)
value = session.scalar(text("SELECT value FROM thread_probe"))
connection_id = id(connection.connection.dbapi_connection)
return connection_id, value
with ThreadPoolExecutor(max_workers=2) as executor:
futures = [executor.submit(read_value) for _ in range(2)]
results = [future.result() for future in futures]
assert {value for _, value in results} == {42}
assert len({connection_id for connection_id, _ in results}) == 2
Generated
+3 -3
View File
@@ -2710,14 +2710,14 @@ wheels = [
[[package]]
name = "gitpython"
version = "3.1.54"
version = "3.1.52"
source = { registry = "https://pypi.org/simple" }
dependencies = [
{ name = "gitdb" },
]
sdist = { url = "https://files.pythonhosted.org/packages/5e/d5/3da0b92033887033f4c27f2dd109a303c4ca62813c7b3bb2511edb4777de/gitpython-3.1.54.tar.gz", hash = "sha256:53f2085e24a2cda300eed7c3fc5f1559ae289634b725e98acaf4791940247aa0", size = 225076, upload-time = "2026-07-22T04:08:51.403Z" }
sdist = { url = "https://files.pythonhosted.org/packages/e5/fd/df0bafa4eb5ea2f51e1adee9f7a94c8e62c5d180e65117045dfca3439c8a/gitpython-3.1.52.tar.gz", hash = "sha256:de0a8ad86274c6e75ae8b37dd055ba68f19818c813108642263227b20775b48e", size = 223726, upload-time = "2026-07-16T03:15:59.599Z" }
wheels = [
{ url = "https://files.pythonhosted.org/packages/d1/b9/876f442a28df5c068ca69b0122d5c35e65fd2d2fa9992ea5cb5944ea00a6/gitpython-3.1.54-py3-none-any.whl", hash = "sha256:b90d7b3d9bc0238681d24369130826f0dcdb0ceaa45db67cf1d4ffa4c302dedf", size = 216575, upload-time = "2026-07-22T04:08:50.05Z" },
{ url = "https://files.pythonhosted.org/packages/8d/90/04dff7c1e176bb1c3011ef1647393d368790da710d8dde1cdcfad301f45a/gitpython-3.1.52-py3-none-any.whl", hash = "sha256:79a36ee1f83523214a3f72d56cf1c4e490d577dc61af77e43dfe5862bd9da01a", size = 215366, upload-time = "2026-07-16T03:15:58.239Z" },
]
[[package]]

Some files were not shown because too many files have changed in this diff Show More