Compare commits

..
Author SHA1 Message Date
Jyong d74499c09d update github action 2026-07-22 06:07:12 -04:00
Jyong 29858d0db0 Expand application capabilities and refactor core workflows 2026-07-22 05:09:09 -04:00
Jyong 77732c7bb1 docs: decouple KnowledgeFS from Dify datasets 2026-07-20 06:02:06 -04:00
Jyong dde79fc9bf docs: add Dify KnowledgeFS integration plan 2026-07-20 05:33:56 -04:00
Jyong 4ee43b8afc chore: migrate knowledge-fs source tree
Import the committed KnowledgeFS snapshot dc4072ee302317145612087ce7440851dc329fd0 under knowledge-fs/ without its Git history, local IDE settings, or build artifacts.
2026-07-20 04:54:20 -04:00
2755 changed files with 571310 additions and 49829 deletions
@@ -32,11 +32,12 @@ Keep this skill focused on Cucumber, Playwright, and package-level E2E guidance.
- `e2e/` uses Cucumber for scenarios and Playwright as the browser layer.
- `DifyWorld` is the per-scenario context object. Type `this` as `DifyWorld` and use `async function`, not arrow functions.
- Keep glue organized by capability under `e2e/features/step-definitions/`; use `common/` only for broadly reusable steps.
- Treat `e2e/AGENTS.md`, `features/support/hooks.ts`, and the Cucumber configuration as the owners of current session and tag semantics. Verify them when behavior depends on session state instead of copying a tag inventory into this skill.
- Browser session behavior comes from `features/support/hooks.ts`:
- default: authenticated session with shared storage state
- `@unauthenticated`: clean browser context
- `@authenticated`: readability/selective-run tag only unless implementation changes
- `@fresh`: only for `e2e:full*` flows
- Do not import Playwright Test runner patterns that bypass the current Cucumber + `DifyWorld` architecture unless the task is explicitly about changing that architecture.
- Perform the behavior under test through Playwright. APIs are allowed for setup, seed preparation, persistence polling, and cleanup, but ordinary Console JSON and representable multipart operations must use the scenario- or process-owned generated oRPC client with request and response validation enabled. Keep the setup/cleanup API identity independent from an unauthenticated or logged-out behavior browser.
- Consume generated operations directly. Do not add one-to-one API wrappers, handwritten endpoint URLs, response DTO casts, duplicate schemas, global mutable clients, or TanStack Query caching in Cucumber. Keep helpers only for real fixture construction, multi-operation orchestration, invariants, polling, derived test views, or protocol adapters.
- Keep SSE, binary, redirect-only, external-service, and readiness exceptions centralized under their protocol owner. A contract mismatch must fail and be fixed at the backend schema owner followed by regeneration; never weaken validation to make E2E pass.
## Workflow
@@ -65,7 +66,7 @@ Keep this skill focused on Cucumber, Playwright, and package-level E2E guidance.
- If a product element has real user-facing semantics but no accessible name, prefer fixing that accessible contract over adding a test id.
5. Validate narrowly.
- Run the narrowest tagged scenario or flow that exercises the change.
- Run the package-required static checks documented in `e2e/AGENTS.md`.
- Run `vpr lint --fix --quiet` from the repository root and `pnpm -C e2e type-check`.
- Broaden verification only when the change affects hooks, tags, setup, or shared step semantics.
## Review Checklist
@@ -76,8 +77,6 @@ Keep this skill focused on Cucumber, Playwright, and package-level E2E guidance.
- Are locators user-facing and assertions web-first?
- Does the change introduce hidden coupling across scenarios, tags, or instance state?
- Does it document or implement behavior that differs from the real hooks or configuration?
- Does setup/cleanup use the generated client directly, with any remaining helper owning more than a one-to-one endpoint forward?
- Is every raw HTTP call a documented protocol or infrastructure exception rather than an ordinary Console operation?
Lead findings with correctness, flake risk, and architecture drift.
+15
View File
@@ -0,0 +1,15 @@
{
"hooks": {
"PreToolUse": [
{
"matcher": "Bash",
"hooks": [
{
"type": "command",
"command": "npx -y block-no-verify@1.1.1"
}
]
}
]
}
}
+1
View File
@@ -8,6 +8,7 @@
**/*.pyc
**/.mypy_cache
**/.ruff_cache
knowledge-fs/
.git
.github
*.md
+44 -17
View File
@@ -8,6 +8,7 @@
# Lint bulk suppression baselines.
/oxlint-suppressions.json
/eslint-suppressions.json
# CODEOWNERS file
/.github/CODEOWNERS @laipz8200 @crazywoola
@@ -32,9 +33,31 @@
# Backend (default owner, more specific rules below will override)
/api/ @QuantumGhost
# Backend - MCP
/api/core/mcp/ @Nov1c444
/api/core/entities/mcp_provider.py @Nov1c444
/api/services/tools/mcp_tools_manage_service.py @Nov1c444
/api/controllers/mcp/ @Nov1c444
/api/controllers/console/app/mcp_server.py @Nov1c444
# Backend - Tests
/api/tests/ @laipz8200 @QuantumGhost
/api/tests/**/*mcp* @Nov1c444
# Backend - Workflow - Engine (Core graph execution engine)
/api/core/workflow/graph_engine/ @laipz8200 @QuantumGhost
/api/core/workflow/runtime/ @laipz8200 @QuantumGhost
/api/core/workflow/graph/ @laipz8200 @QuantumGhost
/api/core/workflow/graph_events/ @laipz8200 @QuantumGhost
/api/core/workflow/node_events/ @laipz8200 @QuantumGhost
# Backend - Workflow - Nodes (Agent, Iteration, Loop, LLM)
/api/core/workflow/nodes/agent/ @Nov1c444
/api/core/workflow/nodes/iteration/ @Nov1c444
/api/core/workflow/nodes/loop/ @Nov1c444
/api/core/workflow/nodes/llm/ @Nov1c444
# Backend - RAG (Retrieval Augmented Generation)
/api/core/rag/ @JohnJyong
/api/services/rag_pipeline/ @JohnJyong
@@ -88,6 +111,7 @@
/api/core/app/layers/trigger_post_layer.py @CourTeous33
/api/services/trigger/ @CourTeous33
/api/models/trigger.py @CourTeous33
/api/fields/workflow_trigger_fields.py @CourTeous33
/api/repositories/workflow_trigger_log_repository.py @CourTeous33
/api/repositories/sqlalchemy_workflow_trigger_log_repository.py @CourTeous33
/api/libs/schedule_utils.py @CourTeous33
@@ -112,11 +136,11 @@
/api/controllers/console/billing/ @hj24 @zyssyz123
# Backend - Enterprise
/api/configs/enterprise/ @GareArc
/api/services/enterprise/ @GareArc
/api/services/feature_service.py @GareArc
/api/controllers/console/feature.py @GareArc
/api/controllers/web/feature.py @GareArc
/api/configs/enterprise/ @GarfieldDai @GareArc
/api/services/enterprise/ @GarfieldDai @GareArc
/api/services/feature_service.py @GarfieldDai @GareArc
/api/controllers/console/feature.py @GarfieldDai @GareArc
/api/controllers/web/feature.py @GarfieldDai @GareArc
# Backend - Database Migrations
/api/migrations/ @snakevash @laipz8200 @MRZHUH
@@ -129,6 +153,7 @@
# Frontend - Platform and Features
/web/config/ @lyzno1
/web/contract/ @lyzno1
/web/env.ts @lyzno1
/web/features/ @lyzno1
/web/hooks/ @lyzno1
@@ -187,6 +212,7 @@
/web/app/components/rag-pipeline/store/ @iamjoel @zxhlyh
# Frontend - RAG - Documents List
/web/app/components/datasets/documents/list.tsx @iamjoel @WTW0313
/web/app/components/datasets/documents/create-from-pipeline/ @iamjoel @WTW0313
# Frontend - RAG - Segments List
@@ -205,22 +231,22 @@
/web/app/components/plugins/marketplace/ @iamjoel @Yessenia-d
# Frontend - Login and Registration
/web/app/signin/ @iamjoel
/web/app/signup/ @iamjoel
/web/app/reset-password/ @iamjoel
/web/app/install/ @iamjoel
/web/app/init/ @iamjoel
/web/app/forgot-password/ @iamjoel
/web/app/account/ @iamjoel
/web/app/signin/ @douxc @iamjoel
/web/app/signup/ @douxc @iamjoel
/web/app/reset-password/ @douxc @iamjoel
/web/app/install/ @douxc @iamjoel
/web/app/init/ @douxc @iamjoel
/web/app/forgot-password/ @douxc @iamjoel
/web/app/account/ @douxc @iamjoel
# Frontend - Service Authentication
/web/service/base.ts @iamjoel
/web/service/base.ts @douxc @iamjoel
# Frontend - WebApp Authentication and Access Control
/web/app/(shareLayout)/components/ @iamjoel
/web/app/(shareLayout)/webapp-signin/ @iamjoel
/web/app/(shareLayout)/webapp-reset-password/ @iamjoel
/web/app/components/app/app-access-control/ @iamjoel
/web/app/(shareLayout)/components/ @douxc @iamjoel
/web/app/(shareLayout)/webapp-signin/ @douxc @iamjoel
/web/app/(shareLayout)/webapp-reset-password/ @douxc @iamjoel
/web/app/components/app/app-access-control/ @douxc @iamjoel
# Frontend - Explore Page
/web/app/components/explore/ @CodingOnStar @iamjoel
@@ -239,6 +265,7 @@
/web/app/components/base/**/*.spec.tsx @hyoban @CodingOnStar
# Frontend - Utils and Hooks
/web/utils/classnames.ts @iamjoel @zxhlyh
/web/utils/time.ts @iamjoel @zxhlyh
/web/utils/format.ts @iamjoel @zxhlyh
/web/utils/clipboard.ts @iamjoel @zxhlyh
+9
View File
@@ -1,6 +1,15 @@
version: 2
updates:
- package-ecosystem: "npm"
directory: "/knowledge-fs"
open-pull-requests-limit: 10
schedule:
interval: "weekly"
groups:
knowledge-fs-dependencies:
patterns:
- "*"
- package-ecosystem: "uv"
directory: "/api"
open-pull-requests-limit: 10
+6 -3
View File
@@ -16,7 +16,7 @@ concurrency:
jobs:
api-unit:
name: API Unit Tests
runs-on: depot-ubuntu-24.04
runs-on: depot-ubuntu-24.04-4
env:
COVERAGE_FILE: coverage-unit
defaults:
@@ -47,6 +47,9 @@ jobs:
- name: Install dependencies
run: uv sync --project api --dev
- name: Run dify config tests
run: uv run --project api pytest api/tests/unit_tests/configs/test_env_consistency.py
- name: Run Unit Tests
run: |
uv run --project api pytest \
@@ -72,7 +75,7 @@ jobs:
api-integration:
name: API Integration Tests
runs-on: depot-ubuntu-24.04
runs-on: depot-ubuntu-24.04-4
env:
COVERAGE_FILE: coverage-integration
STORAGE_TYPE: opendal
@@ -126,7 +129,7 @@ jobs:
api-coverage:
name: API Coverage
runs-on: depot-ubuntu-24.04
runs-on: depot-ubuntu-24.04-4
needs:
- api-unit
- api-integration
+1 -1
View File
@@ -173,7 +173,7 @@ jobs:
create-manifest:
needs: build
runs-on: depot-ubuntu-24.04
runs-on: depot-ubuntu-24.04-4
if: github.repository == 'langgenius/dify'
strategy:
matrix:
+2 -2
View File
@@ -23,7 +23,7 @@ concurrency:
jobs:
validate:
name: validate manifest + resolve target Dify release
runs-on: depot-ubuntu-24.04
runs-on: depot-ubuntu-24.04-4
if: github.repository == 'langgenius/dify'
permissions:
contents: read
@@ -87,7 +87,7 @@ jobs:
release:
name: build + attach standalone binaries (all targets)
needs: validate
runs-on: depot-ubuntu-24.04
runs-on: depot-ubuntu-24.04-4
permissions:
contents: write
defaults:
+2 -2
View File
@@ -9,7 +9,7 @@ concurrency:
jobs:
db-migration-test-postgres:
runs-on: depot-ubuntu-24.04
runs-on: depot-ubuntu-24.04-4
steps:
- name: Checkout code
@@ -59,7 +59,7 @@ jobs:
run: uv run --directory api flask upgrade-db
db-migration-test-mysql:
runs-on: depot-ubuntu-24.04
runs-on: depot-ubuntu-24.04-4
steps:
- name: Checkout code
+1 -1
View File
@@ -13,7 +13,7 @@ on:
jobs:
deploy:
runs-on: depot-ubuntu-24.04
runs-on: depot-ubuntu-24.04-4
if: |
github.event.workflow_run.conclusion == 'success' &&
github.event.workflow_run.head_branch == 'deploy/agent'
+1 -1
View File
@@ -10,7 +10,7 @@ on:
jobs:
deploy:
runs-on: depot-ubuntu-24.04
runs-on: depot-ubuntu-24.04-4
if: |
github.event.workflow_run.conclusion == 'success' &&
github.event.workflow_run.head_branch == 'deploy/dev'
+1 -1
View File
@@ -13,7 +13,7 @@ on:
jobs:
deploy:
runs-on: depot-ubuntu-24.04
runs-on: depot-ubuntu-24.04-4
if: |
github.event.workflow_run.conclusion == 'success' &&
github.event.workflow_run.head_branch == 'deploy/enterprise'
+1 -1
View File
@@ -13,7 +13,7 @@ on:
jobs:
deploy:
runs-on: depot-ubuntu-24.04
runs-on: depot-ubuntu-24.04-4
if: |
github.event.workflow_run.conclusion == 'success' &&
github.event.workflow_run.head_branch == 'deploy/saas'
+1 -1
View File
@@ -22,7 +22,7 @@ concurrency:
jobs:
check-cherry-pick-provenance:
name: Require cherry-pick provenance
runs-on: depot-ubuntu-24.04
runs-on: depot-ubuntu-24.04-4
steps:
- uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0
with:
+504
View File
@@ -0,0 +1,504 @@
name: KnowledgeFS CI
on:
pull_request:
branches: ["main"]
merge_group:
branches: ["main"]
types: [checks_requested]
push:
branches: ["main"]
workflow_dispatch:
permissions:
contents: read
pull-requests: read
concurrency:
group: knowledge-fs-${{ github.workflow }}-${{ github.event.pull_request.number || github.ref }}
cancel-in-progress: true
env:
CI: true
DIFY_KNOWLEDGE_FS_API_IMAGE_NAME: >-
${{ vars.DIFY_KNOWLEDGE_FS_API_IMAGE_NAME || 'langgenius/dify-knowledge-fs-api' }}
jobs:
check-changes:
name: Check KnowledgeFS changes
runs-on: depot-ubuntu-24.04-4
outputs:
knowledge-fs: ${{ steps.changes.outputs.knowledge-fs }}
steps:
- name: Checkout code
uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0
with:
fetch-depth: 0
persist-credentials: false
- name: Detect KnowledgeFS changes
id: changes
uses: dorny/paths-filter@7b450fff21473bca461d4b92ce414b9d0420d706 # v4.0.2
with:
filters: |
knowledge-fs:
- 'knowledge-fs/**'
- 'knowledge-fs/packages/api/src/dify-capability-v2.ts'
- 'knowledge-fs/packages/api/src/knowledge-space-routes.ts'
- 'knowledge-fs/packages/api/src/upload-session-routes.ts'
- 'knowledge-fs/scripts/export-capability-v2-operations.mjs'
- 'knowledge-fs/scripts/export-openapi.mjs'
- 'api/dev/generate_knowledge_fs_contract.py'
- 'api/dev/knowledge_fs_product_contract.py'
- 'api/knowledge-fs-contract.lock.json'
- 'api/knowledge-fs-product-operation-gaps.json'
- 'api/knowledge-fs-product-operations.json'
- 'api/**/knowledge_fs/**'
- 'api/**/*knowledge_fs*'
- 'api/**/*knowledge-fs*'
- 'api/.env.example'
- 'api/app_factory.py'
- 'api/commands/__init__.py'
- 'api/controllers/console/__init__.py'
- 'api/controllers/console/workspace/rbac.py'
- 'api/controllers/service_api/__init__.py'
- 'api/core/agent/base_agent_runner.py'
- 'api/core/app/apps/agent_app/runtime_request_builder.py'
- 'api/core/rbac/entities.py'
- 'api/core/tools/__base/tool_runtime.py'
- 'api/core/tools/builtin_tool/_position.yaml'
- 'api/core/workflow/node_runtime.py'
- 'api/core/workflow/nodes/agent_v2/runtime_request_builder.py'
- 'api/extensions/ext_celery.py'
- 'api/extensions/ext_commands.py'
- 'api/models/__init__.py'
- 'api/services/account_service.py'
- 'api/services/agent_tool_inner_service.py'
- 'api/services/enterprise/rbac_service.py'
- 'api/services/entities/agent_tool_inner.py'
- 'api/services/knowledge_fs/**'
- 'api/services/knowledge_fs_capability.py'
- 'api/tests/unit_tests/dev/test_generate_knowledge_fs_contract.py'
- 'api/tests/unit_tests/controllers/console/workspace/test_rbac.py'
- 'api/tests/unit_tests/core/agent/test_base_agent_runner.py'
- 'api/tests/unit_tests/core/app/apps/agent_app/test_runtime_request_builder.py'
- 'api/tests/unit_tests/core/workflow/nodes/agent_v2/test_runtime_request_builder.py'
- 'api/tests/unit_tests/core/workflow/nodes/tool/test_tool_node_runtime.py'
- 'api/tests/unit_tests/core/workflow/test_node_runtime.py'
- 'api/tests/unit_tests/services/enterprise/test_rbac_service.py'
- 'api/tests/unit_tests/services/test_account_service.py'
- 'api/tests/unit_tests/services/test_agent_tool_inner_service.py'
- 'api/tests/unit_tests/services/test_knowledge_fs_capability.py'
- 'api/tests/unit_tests/services/test_knowledge_fs_product_operations.py'
- 'api/pyproject.toml'
- 'api/uv.lock'
- 'dify-agent/src/dify_agent/layers/dify_core_tools/client.py'
- 'dify-agent/tests/local/dify_agent/layers/dify_core_tools/test_client.py'
- 'packages/contracts/generated/api/console/**'
- 'packages/contracts/generated/api/service/**'
- 'docker/.env.example'
- 'docker/README.md'
- 'docker/dify-env-sync.py'
- 'docker/dify-env-sync.sh'
- 'docker/docker-compose-template.yaml'
- 'docker/docker-compose.yaml'
- 'docker/envs/core-services/api.env.example'
- 'docker/envs/core-services/knowledge-fs.env.example'
- 'docker/generate_docker_compose'
- 'docs/design/knowledge-fs*'
- '.github/dependabot.yml'
- '.github/workflows/knowledge-fs-ci.yml'
build:
name: Build KnowledgeFS API production image
needs: check-changes
if: needs.check-changes.outputs.knowledge-fs == 'true' || github.event_name == 'workflow_dispatch'
runs-on: depot-ubuntu-24.04-4
steps:
- name: Checkout code
uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0
with:
persist-credentials: false
- name: Set up Docker Buildx
uses: docker/setup-buildx-action@bb05f3f5519dd87d3ba754cc423b652a5edd6d2c # v4.2.0
- name: Login to Docker Hub
if: github.event_name == 'workflow_dispatch' || (github.event_name == 'push' && github.ref == 'refs/heads/main')
uses: docker/login-action@af1e73f918a031802d376d3c8bbc3fe56130a9b0 # v4.4.0
with:
username: ${{ secrets.DOCKERHUB_USER }}
password: ${{ secrets.DOCKERHUB_TOKEN }}
- name: Extract KnowledgeFS image metadata
id: meta
uses: docker/metadata-action@dc802804100637a589fabce1cb79ff13a1411302 # v6.2.0
with:
images: ${{ env.DIFY_KNOWLEDGE_FS_API_IMAGE_NAME }}
tags: |
type=raw,value=latest,enable=${{ github.ref == 'refs/heads/main' }}
type=ref,event=branch
type=sha,format=long
- name: Build KnowledgeFS API image
uses: docker/build-push-action@53b7df96c91f9c12dcc8a07bcb9ccacbed38856a # v7.3.0
with:
context: ./knowledge-fs
file: ./knowledge-fs/apps/api/Dockerfile
labels: ${{ steps.meta.outputs.labels }}
platforms: linux/amd64
push: ${{ github.event_name == 'workflow_dispatch' || (github.event_name == 'push' && github.ref == 'refs/heads/main') }}
tags: ${{ steps.meta.outputs.tags }}
quality:
name: Run KnowledgeFS quality and contract gates
needs: check-changes
if: needs.check-changes.outputs.knowledge-fs == 'true' || github.event_name == 'workflow_dispatch'
runs-on: depot-ubuntu-24.04-4
defaults:
run:
shell: bash
working-directory: ./knowledge-fs
steps:
- name: Checkout code
uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0
with:
fetch-depth: 0
persist-credentials: false
- name: Setup pnpm
uses: pnpm/action-setup@0ebf47130e4866e96fce0953f49152a61190b271 # v6.0.9
with:
package_json_file: knowledge-fs/package.json
run_install: false
- name: Setup Node
uses: actions/setup-node@820762786026740c76f36085b0efc47a31fe5020 # v7.0.0
with:
node-version: 22
cache: pnpm
cache-dependency-path: knowledge-fs/pnpm-lock.yaml
- name: Install KnowledgeFS dependencies
run: pnpm install --frozen-lockfile
- name: Scan KnowledgeFS secrets
run: pnpm security:secrets
- name: Audit KnowledgeFS production dependencies
run: pnpm security:dependencies
- name: Run KnowledgeFS checks
run: pnpm check
- name: Build KnowledgeFS
run: pnpm build
- name: Lint KnowledgeFS
run: pnpm lint
- name: Setup UV and Python
uses: astral-sh/setup-uv@11f9893b081a58869d3b5fccaea48c9e9e46f990 # v8.3.2
with:
enable-cache: true
python-version: "3.12"
cache-dependency-glob: |
api/uv.lock
dify-agent/uv.lock
- name: Verify Dify dependency lock
working-directory: .
run: uv lock --project api --check
- name: Install Dify contract dependencies
working-directory: .
run: uv sync --project api --locked --dev
- name: Collect Dify KnowledgeFS gate targets
working-directory: .
run: |
set -euo pipefail
target_dir="${RUNNER_TEMP:?}/knowledge-fs-ci-targets"
mkdir -p "$target_dir"
production_targets=()
add_production_target() {
local path="$1"
if [[ ! -f "$path" ]]; then
echo "required Dify KnowledgeFS production target is missing: $path" >&2
exit 1
fi
production_targets+=("$path")
}
while IFS= read -r -d '' path; do
if [[ "$path" == *knowledge_fs* ]]; then
add_production_target "$path"
fi
done < <(
find api \
\( -path 'api/.venv' -o -path 'api/tests' -o -path 'api/storage' \) -prune \
-o -type f -name '*.py' -print0
)
production_touchpoints=(
api/app_factory.py
api/commands/__init__.py
api/controllers/console/__init__.py
api/controllers/console/workspace/rbac.py
api/controllers/service_api/__init__.py
api/core/agent/base_agent_runner.py
api/core/app/apps/agent_app/runtime_request_builder.py
api/core/rbac/entities.py
api/core/tools/__base/tool_runtime.py
api/core/workflow/node_runtime.py
api/core/workflow/nodes/agent_v2/runtime_request_builder.py
api/extensions/ext_celery.py
api/extensions/ext_commands.py
api/models/__init__.py
api/services/account_service.py
api/services/agent_tool_inner_service.py
api/services/enterprise/rbac_service.py
api/services/entities/agent_tool_inner.py
)
for path in "${production_touchpoints[@]}"; do
add_production_target "$path"
done
printf '%s\0' "${production_touchpoints[@]}" > "$target_dir/glue-files"
if ((${#production_targets[@]} == 0)); then
echo "Dify KnowledgeFS production target set is empty" >&2
exit 1
fi
printf '%s\0' "${production_targets[@]}" > "$target_dir/production-files"
test_targets=()
add_test_target() {
local path="$1"
if [[ ! -f "$path" ]]; then
echo "required Dify KnowledgeFS unit test is missing: $path" >&2
exit 1
fi
test_targets+=("$path")
}
while IFS= read -r -d '' path; do
if [[ "$path" == *knowledge_fs* ]]; then
add_test_target "$path"
fi
done < <(find api/tests/unit_tests -type f -name '*.py' -print0)
test_touchpoints=(
api/tests/unit_tests/controllers/console/workspace/test_rbac.py
api/tests/unit_tests/core/agent/test_base_agent_runner.py
api/tests/unit_tests/core/app/apps/agent_app/test_runtime_request_builder.py
api/tests/unit_tests/core/workflow/nodes/agent_v2/test_runtime_request_builder.py
api/tests/unit_tests/core/workflow/nodes/tool/test_tool_node_runtime.py
api/tests/unit_tests/core/workflow/test_node_runtime.py
api/tests/unit_tests/services/enterprise/test_rbac_service.py
api/tests/unit_tests/services/test_account_service.py
api/tests/unit_tests/services/test_agent_tool_inner_service.py
)
for path in "${test_touchpoints[@]}"; do
add_test_target "$path"
done
required_test_scopes=(
/commands/
/configs/
/controllers/
/core/agent/
/core/app/apps/agent_app/
/core/tools/builtin_tool/providers/knowledge_fs/
/core/workflow/
/dev/
/extensions/
/migrations/
/models/
/repositories/
/services/
/tasks/
)
for required_scope in "${required_test_scopes[@]}"; do
scope_found=false
for path in "${test_targets[@]}"; do
if [[ "$path" == *"$required_scope"* ]]; then
scope_found=true
break
fi
done
if [[ "$scope_found" != true ]]; then
echo "required Dify KnowledgeFS test scope is empty: $required_scope" >&2
exit 1
fi
done
if ((${#test_targets[@]} == 0)); then
echo "Dify KnowledgeFS unit test target set is empty" >&2
exit 1
fi
printf '%s\0' "${test_targets[@]}" > "$target_dir/unit-test-files"
- name: Lint Dify KnowledgeFS integration
working-directory: .
run: |
set -euo pipefail
targets=()
while IFS= read -r -d '' path; do
targets+=("$path")
done < "${RUNNER_TEMP:?}/knowledge-fs-ci-targets/production-files"
if ((${#targets[@]} == 0)); then
echo "Dify KnowledgeFS production target manifest is empty" >&2
exit 1
fi
uv run --project api --dev ruff format --check "${targets[@]}"
uv run --project api --dev ruff check "${targets[@]}"
- name: Type-check Dify KnowledgeFS integration
working-directory: .
run: |
set -euo pipefail
targets=()
while IFS= read -r -d '' path; do
targets+=("$path")
done < "${RUNNER_TEMP:?}/knowledge-fs-ci-targets/production-files"
if ((${#targets[@]} == 0)); then
echo "Dify KnowledgeFS production target manifest is empty" >&2
exit 1
fi
PYREFLY_OUTPUT_FORMAT=github ./dev/pyrefly-check-local "${targets[@]}"
mypy_targets=()
for path in "${targets[@]}"; do
if [[ "$path" != api/migrations/* ]]; then
mypy_targets+=("${path#api/}")
fi
done
if ((${#mypy_targets[@]} == 0)); then
echo "Dify KnowledgeFS Mypy target set is empty" >&2
exit 1
fi
uv run --directory api --dev mypy \
--explicit-package-bases \
--exclude-gitignore \
--exclude '(^|/)conftest\.py$' \
--exclude 'tests/' \
--exclude 'migrations/' \
--check-untyped-defs \
--disable-error-code=import-untyped \
"${mypy_targets[@]}"
- name: Test Dify KnowledgeFS unit surface
working-directory: .
env:
COVERAGE_FILE: ${{ runner.temp }}/dify-knowledge-fs.coverage
run: |
set -euo pipefail
targets=()
while IFS= read -r -d '' path; do
targets+=("$path")
done < "${RUNNER_TEMP:?}/knowledge-fs-ci-targets/unit-test-files"
if ((${#targets[@]} == 0)); then
echo "Dify KnowledgeFS unit test target manifest is empty" >&2
exit 1
fi
uv run --project api --dev coverage run --branch --source=api -m pytest "${targets[@]}" --no-cov -q
- name: Enforce Dify KnowledgeFS focused coverage
working-directory: .
env:
COVERAGE_FILE: ${{ runner.temp }}/dify-knowledge-fs.coverage
KNOWLEDGE_FS_COVERAGE_BASE: >-
${{ github.event.pull_request.base.sha || github.event.merge_group.base_sha || github.event.before || '' }}
run: |
set -euo pipefail
report="${RUNNER_TEMP:?}/dify-knowledge-fs-coverage.json"
uv run --project api --dev coverage json --show-contexts -o "$report"
uv run --project api --dev python api/dev/check_knowledge_fs_coverage.py \
--coverage-json "$report" \
--glue-manifest "${RUNNER_TEMP:?}/knowledge-fs-ci-targets/glue-files" \
--base "$KNOWLEDGE_FS_COVERAGE_BASE" \
--minimum 90 \
--glue-minimum 90
- name: Verify Dify KnowledgeFS contract
working-directory: .
run: uv run --project api python api/dev/generate_knowledge_fs_contract.py --check
- name: Verify Dify Agent dependency lock
working-directory: .
run: uv lock --project dify-agent --check
- name: Install Dify Agent gate dependencies
working-directory: .
run: uv sync --project dify-agent --locked --dev
- name: Lint Dify Agent KnowledgeFS integration
working-directory: ./dify-agent
run: |
uv run --project . --dev ruff format --check \
src/dify_agent/layers/dify_core_tools/client.py \
tests/local/dify_agent/layers/dify_core_tools/test_client.py
uv run --project . --dev ruff check \
src/dify_agent/layers/dify_core_tools/client.py \
tests/local/dify_agent/layers/dify_core_tools/test_client.py
- name: Type-check Dify Agent KnowledgeFS integration
working-directory: ./dify-agent
run: >-
uv run --project . --dev basedpyright --level error
src/dify_agent/layers/dify_core_tools/client.py
tests/local/dify_agent/layers/dify_core_tools/test_client.py
- name: Test Dify Agent KnowledgeFS integration
working-directory: ./dify-agent
run: >-
uv run --project . --dev python -m pytest
tests/local/dify_agent/layers/dify_core_tools/test_client.py
-q
skip:
name: Skip KnowledgeFS quality and contract gates
needs: check-changes
if: needs.check-changes.outputs.knowledge-fs != 'true' && github.event_name != 'workflow_dispatch'
runs-on: depot-ubuntu-24.04-4
steps:
- name: Report skipped KnowledgeFS checks
run: echo "No KnowledgeFS-related changes detected; skipping KnowledgeFS checks."
final:
name: KnowledgeFS CI
if: ${{ always() }}
needs:
- check-changes
- build
- quality
- skip
runs-on: depot-ubuntu-24.04-4
steps:
- name: Finalize KnowledgeFS CI status
env:
EVENT_NAME: ${{ github.event_name }}
BUILD_RESULT: ${{ needs.build.result }}
KNOWLEDGE_FS_CHANGED: ${{ needs.check-changes.outputs.knowledge-fs }}
QUALITY_RESULT: ${{ needs.quality.result }}
SKIP_RESULT: ${{ needs.skip.result }}
run: |
if [[ "$EVENT_NAME" == 'workflow_dispatch' || "$KNOWLEDGE_FS_CHANGED" == 'true' ]]; then
if [[ "$BUILD_RESULT" == 'success' && "$QUALITY_RESULT" == 'success' ]]; then
echo "KnowledgeFS build and checks ran successfully."
exit 0
fi
echo "KnowledgeFS build or checks failed: build=$BUILD_RESULT quality=$QUALITY_RESULT" >&2
exit 1
fi
if [[ "$SKIP_RESULT" == 'success' ]]; then
echo "KnowledgeFS checks were skipped because no related files changed."
exit 0
fi
echo "KnowledgeFS change detection or skip reporting failed with result: $SKIP_RESULT" >&2
exit 1
+1 -1
View File
@@ -7,7 +7,7 @@ jobs:
permissions:
contents: read
pull-requests: write
runs-on: depot-ubuntu-24.04
runs-on: depot-ubuntu-24.04-4
steps:
- uses: actions/labeler@b8dd2d9be0f68b860e7dae5dae7d772984eacd6d # v6.2.0
with:
+17 -18
View File
@@ -21,7 +21,7 @@ concurrency:
jobs:
pre_job:
name: Skip Duplicate Checks
runs-on: depot-ubuntu-24.04
runs-on: depot-ubuntu-24.04-4
outputs:
should_skip: ${{ steps.skip_check.outputs.should_skip || 'false' }}
steps:
@@ -37,7 +37,7 @@ jobs:
name: Check Changed Files
needs: pre_job
if: needs.pre_job.outputs.should_skip != 'true'
runs-on: depot-ubuntu-24.04
runs-on: depot-ubuntu-24.04-4
outputs:
api-changed: ${{ steps.changes.outputs.api }}
cli-changed: ${{ steps.changes.outputs.cli }}
@@ -81,6 +81,7 @@ jobs:
- '.npmrc'
- '.nvmrc'
- '.github/workflows/cli-tests.yml'
- '.github/workflows/cli-docker-build.yml'
- '.github/actions/setup-web/**'
web:
- 'web/**'
@@ -163,7 +164,7 @@ jobs:
- pre_job
- check-changes
if: needs.pre_job.outputs.should_skip != 'true' && needs.check-changes.outputs.api-changed != 'true'
runs-on: depot-ubuntu-24.04
runs-on: depot-ubuntu-24.04-4
steps:
- name: Report skipped API tests
run: echo "No API-related changes detected; skipping API tests."
@@ -176,7 +177,7 @@ jobs:
- check-changes
- api-tests-run
- api-tests-skip
runs-on: depot-ubuntu-24.04
runs-on: depot-ubuntu-24.04-4
steps:
- name: Finalize API Tests status
env:
@@ -223,7 +224,7 @@ jobs:
- pre_job
- check-changes
if: needs.pre_job.outputs.should_skip != 'true' && needs.check-changes.outputs.cli-changed != 'true'
runs-on: depot-ubuntu-24.04
runs-on: depot-ubuntu-24.04-4
steps:
- name: Report skipped CLI tests
run: echo "No CLI-related changes detected; skipping CLI tests."
@@ -236,7 +237,7 @@ jobs:
- check-changes
- cli-tests-run
- cli-tests-skip
runs-on: depot-ubuntu-24.04
runs-on: depot-ubuntu-24.04-4
steps:
- name: Finalize CLI Tests status
env:
@@ -283,7 +284,7 @@ jobs:
- pre_job
- check-changes
if: needs.pre_job.outputs.should_skip != 'true' && needs.check-changes.outputs.web-changed != 'true'
runs-on: depot-ubuntu-24.04
runs-on: depot-ubuntu-24.04-4
steps:
- name: Report skipped web tests
run: echo "No web-related changes detected; skipping web tests."
@@ -296,7 +297,7 @@ jobs:
- check-changes
- web-tests-run
- web-tests-skip
runs-on: depot-ubuntu-24.04
runs-on: depot-ubuntu-24.04-4
steps:
- name: Finalize Web Tests status
env:
@@ -335,8 +336,6 @@ jobs:
- check-changes
if: needs.pre_job.outputs.should_skip != 'true' && needs.check-changes.outputs.e2e-changed == 'true'
uses: ./.github/workflows/web-e2e.yml
with:
run-external-runtime: false
secrets: inherit
web-e2e-skip:
@@ -345,7 +344,7 @@ jobs:
- pre_job
- check-changes
if: needs.pre_job.outputs.should_skip != 'true' && needs.check-changes.outputs.e2e-changed != 'true'
runs-on: depot-ubuntu-24.04
runs-on: depot-ubuntu-24.04-4
steps:
- name: Report skipped web full-stack e2e
run: echo "No E2E-related changes detected; skipping web full-stack E2E."
@@ -358,7 +357,7 @@ jobs:
- check-changes
- web-e2e-run
- web-e2e-skip
runs-on: depot-ubuntu-24.04
runs-on: depot-ubuntu-24.04-4
steps:
- name: Finalize Web Full-Stack E2E status
env:
@@ -412,7 +411,7 @@ jobs:
- pre_job
- check-changes
if: needs.pre_job.outputs.should_skip != 'true' && needs.check-changes.outputs.vdb-changed != 'true'
runs-on: depot-ubuntu-24.04
runs-on: depot-ubuntu-24.04-4
steps:
- name: Report skipped VDB tests
run: echo "No VDB-related changes detected; skipping VDB tests."
@@ -425,7 +424,7 @@ jobs:
- check-changes
- vdb-tests-run
- vdb-tests-skip
runs-on: depot-ubuntu-24.04
runs-on: depot-ubuntu-24.04-4
steps:
- name: Finalize VDB Tests status
env:
@@ -471,7 +470,7 @@ jobs:
- pre_job
- check-changes
if: needs.pre_job.outputs.should_skip != 'true' && needs.check-changes.outputs.migration-changed != 'true'
runs-on: depot-ubuntu-24.04
runs-on: depot-ubuntu-24.04-4
steps:
- name: Report skipped DB migration tests
run: echo "No migration-related changes detected; skipping DB migration tests."
@@ -484,7 +483,7 @@ jobs:
- check-changes
- db-migration-test-run
- db-migration-test-skip
runs-on: depot-ubuntu-24.04
runs-on: depot-ubuntu-24.04-4
steps:
- name: Finalize DB Migration Test status
env:
@@ -531,7 +530,7 @@ jobs:
- pre_job
- check-changes
if: needs.pre_job.outputs.should_skip != 'true' && needs.check-changes.outputs.sandbox-runtime-changed != 'true'
runs-on: depot-ubuntu-24.04
runs-on: depot-ubuntu-24.04-4
steps:
- name: Report skipped sandbox runtime tests
run: echo "No sandbox-runtime-related changes detected; skipping sandbox runtime tests."
@@ -544,7 +543,7 @@ jobs:
- check-changes
- sandbox-runtime-tests-run
- sandbox-runtime-tests-skip
runs-on: depot-ubuntu-24.04
runs-on: depot-ubuntu-24.04-4
steps:
- name: Finalize Sandbox Runtime Tests status
env:
+1 -17
View File
@@ -14,7 +14,7 @@ concurrency:
jobs:
check-changes:
name: Check Changed Files
runs-on: depot-ubuntu-24.04
runs-on: depot-ubuntu-24.04-4
outputs:
external-e2e-changed: ${{ steps.changes.outputs.external_e2e }}
steps:
@@ -26,30 +26,19 @@ jobs:
external_e2e:
- 'e2e/features/agent-v2/**'
- 'e2e/features/step-definitions/agent-v2/**'
- 'e2e/features/step-definitions/common/**'
- 'e2e/features/support/**'
- 'e2e/fixtures/auth.ts'
- 'e2e/fixtures/test-materials/**'
- 'e2e/scripts/**'
- 'e2e/support/**'
- 'e2e/cucumber.config.ts'
- 'e2e/package.json'
- 'e2e/test-env.ts'
- 'e2e/tsconfig.json'
- 'e2e/tsx-register.js'
- 'package.json'
- 'pnpm-lock.yaml'
- '.nvmrc'
- '.github/workflows/post-merge.yml'
- '.github/workflows/web-e2e.yml'
- '.github/actions/setup-web/**'
- 'docker/docker-compose.middleware.yaml'
- 'docker/envs/middleware.env.example'
- 'dify-agent/**'
- 'dify-agent-runtime/**'
- 'api/pyproject.toml'
- 'api/uv.lock'
- 'api/tests/integration_tests/.env.example'
- 'api/clients/agent_backend/**'
- 'api/core/app/apps/agent_app/**'
- 'api/core/workflow/nodes/agent_v2/**'
@@ -59,13 +48,8 @@ jobs:
- 'api/services/plugin/**'
- 'api/core/tools/**'
- 'api/services/tools/**'
- 'packages/contracts/package.json'
- 'packages/contracts/generated/api/console/agent/**'
- 'packages/contracts/generated/api/console/apps/**'
- 'packages/contracts/generated/api/console/datasets/**'
- 'packages/contracts/generated/api/console/orpc.gen.ts'
- 'packages/contracts/generated/api/console/workspaces/**'
- 'packages/contracts/generated/api/service/**'
- 'web/features/agent-v2/**'
- 'web/app/(commonLayout)/agents/**'
- 'web/app/(commonLayout)/@detailSidebar/agents/**'
+1 -1
View File
@@ -12,7 +12,7 @@ permissions: {}
jobs:
comment:
name: Comment PR with pyrefly diff
runs-on: depot-ubuntu-24.04
runs-on: depot-ubuntu-24.04-4
permissions:
actions: read
contents: read
+1 -1
View File
@@ -10,7 +10,7 @@ permissions:
jobs:
pyrefly-diff:
runs-on: depot-ubuntu-24.04
runs-on: depot-ubuntu-24.04-4
permissions:
contents: read
issues: write
@@ -12,7 +12,7 @@ permissions: {}
jobs:
comment:
name: Comment PR with type coverage
runs-on: depot-ubuntu-24.04
runs-on: depot-ubuntu-24.04-4
permissions:
actions: read
contents: read
+1 -1
View File
@@ -10,7 +10,7 @@ permissions:
jobs:
pyrefly-type-coverage:
runs-on: depot-ubuntu-24.04
runs-on: depot-ubuntu-24.04-4
permissions:
contents: read
issues: write
+3 -3
View File
@@ -13,7 +13,7 @@ concurrency:
jobs:
sandbox-runtime-unit:
name: Sandbox Runtime Unit Tests
runs-on: depot-ubuntu-24.04
runs-on: depot-ubuntu-24.04-4
defaults:
run:
shell: bash
@@ -37,7 +37,7 @@ jobs:
sandbox-runtime-lint:
name: Sandbox Runtime Lint
runs-on: depot-ubuntu-24.04
runs-on: depot-ubuntu-24.04-4
defaults:
run:
shell: bash
@@ -64,7 +64,7 @@ jobs:
sandbox-runtime-integration:
name: Sandbox Runtime Integration Tests
runs-on: depot-ubuntu-24.04
runs-on: depot-ubuntu-24.04-4
defaults:
run:
shell: bash
+1 -1
View File
@@ -16,7 +16,7 @@ jobs:
name: Validate PR title
permissions:
pull-requests: read
runs-on: depot-ubuntu-24.04
runs-on: depot-ubuntu-24.04-4
steps:
- name: Complete merge group check
if: github.event_name == 'merge_group'
+1 -1
View File
@@ -12,7 +12,7 @@ on:
jobs:
stale:
runs-on: depot-ubuntu-24.04
runs-on: depot-ubuntu-24.04-4
permissions:
issues: write
pull-requests: write
+3 -3
View File
@@ -19,7 +19,7 @@ permissions:
jobs:
python-style:
name: Python Style
runs-on: depot-ubuntu-24.04
runs-on: depot-ubuntu-24.04-4
steps:
- name: Checkout code
@@ -83,7 +83,7 @@ jobs:
web-style:
name: Web Style
runs-on: depot-ubuntu-24.04
runs-on: depot-ubuntu-24.04-4
defaults:
run:
working-directory: ./web
@@ -182,7 +182,7 @@ jobs:
superlinter:
name: SuperLinter
runs-on: depot-ubuntu-24.04
runs-on: depot-ubuntu-24.04-4
steps:
- name: Checkout code
+1 -1
View File
@@ -17,7 +17,7 @@ concurrency:
jobs:
build:
name: unit test for Node.js SDK
runs-on: depot-ubuntu-24.04
runs-on: depot-ubuntu-24.04-4
defaults:
run:
+1 -1
View File
@@ -35,7 +35,7 @@ concurrency:
jobs:
translate:
if: github.repository == 'langgenius/dify'
runs-on: depot-ubuntu-24.04
runs-on: depot-ubuntu-24.04-4
timeout-minutes: 120
steps:
+1 -1
View File
@@ -16,7 +16,7 @@ concurrency:
jobs:
trigger:
if: github.repository == 'langgenius/dify'
runs-on: depot-ubuntu-24.04
runs-on: depot-ubuntu-24.04-4
timeout-minutes: 5
steps:
+1 -1
View File
@@ -16,7 +16,7 @@ jobs:
test:
name: Full VDB Tests
if: github.repository == 'langgenius/dify'
runs-on: depot-ubuntu-24.04
runs-on: depot-ubuntu-24.04-4
strategy:
matrix:
python-version:
+1 -1
View File
@@ -13,7 +13,7 @@ concurrency:
jobs:
test:
name: VDB Smoke Tests
runs-on: depot-ubuntu-24.04
runs-on: depot-ubuntu-24.04-4
strategy:
matrix:
python-version:
+14 -7
View File
@@ -4,9 +4,9 @@ on:
workflow_call:
inputs:
run-external-runtime:
description: Run only the prepared and external runtime suite instead of the core suites.
required: true
required: false
type: boolean
default: false
permissions:
contents: read
@@ -46,7 +46,6 @@ jobs:
run: uv sync --project api --dev
- name: Run E2E support unit tests
if: ${{ !inputs.run-external-runtime }}
working-directory: ./e2e
run: vp run test:unit
@@ -55,7 +54,6 @@ jobs:
run: vp run e2e:install
- name: Run isolated source-api and built-web Cucumber E2E tests
if: ${{ !inputs.run-external-runtime }}
working-directory: ./e2e
env:
E2E_ADMIN_EMAIL: e2e-admin@example.com
@@ -66,7 +64,7 @@ jobs:
run: vp run e2e:full
- name: Preserve Chromium E2E report and logs
if: ${{ !cancelled() && !inputs.run-external-runtime }}
if: ${{ !cancelled() }}
run: |
if [[ -d e2e/cucumber-report ]]; then
mv e2e/cucumber-report e2e/cucumber-report-non-external
@@ -76,7 +74,6 @@ jobs:
fi
- name: Run WebKit keyboard and browser smoke tests
if: ${{ !inputs.run-external-runtime }}
working-directory: ./e2e
env:
E2E_ADMIN_EMAIL: e2e-admin@example.com
@@ -102,7 +99,7 @@ jobs:
vp run e2e -- --tags '@browser-smoke'
- name: Preserve WebKit E2E report and logs
if: ${{ !cancelled() && !inputs.run-external-runtime }}
if: ${{ !cancelled() }}
run: |
if [[ -d e2e/cucumber-report ]]; then
mv e2e/cucumber-report e2e/cucumber-report-webkit
@@ -138,6 +135,16 @@ jobs:
exit 1
fi
if [[ -d cucumber-report ]]; then
rm -rf cucumber-report-non-external
mv cucumber-report cucumber-report-non-external
fi
if [[ -d .logs ]]; then
rm -rf .logs-non-external
mv .logs .logs-non-external
fi
teardown_external_runtime() {
local run_status=$?
trap - EXIT
+2
View File
@@ -17,6 +17,8 @@ jobs:
test:
name: Web Tests (${{ matrix.shardIndex }}/${{ matrix.shardTotal }})
runs-on: depot-ubuntu-24.04-4
env:
VITEST_COVERAGE_SCOPE: app-components
strategy:
fail-fast: false
matrix:
+5
View File
@@ -30,6 +30,11 @@ share/python-wheels/
*.egg
MANIFEST
# KnowledgeFS is an independently rooted TypeScript workspace. Its admin `lib`
# directory contains source files rather than Python build output.
!/knowledge-fs/apps/admin/lib/
!/knowledge-fs/apps/admin/lib/**
# PyInstaller
# Usually these files are written by a python script from a template
# before PyInstaller builds the exe, so as to inject date/other infos into it.
-2
View File
@@ -107,7 +107,6 @@ test:
echo "Target: $(TARGET_TESTS)"; \
uv run --project api --dev pytest $(TARGET_TESTS); \
else \
set -e; \
echo "Running backend unit tests"; \
uv run --project api --dev pytest -p no:benchmark --timeout "$${PYTEST_TIMEOUT:-20}" -n auto \
api/tests/unit_tests \
@@ -125,7 +124,6 @@ test-all:
echo "Target: $(TARGET_TESTS)"; \
uv run --project api --dev pytest $(TARGET_TESTS); \
else \
set -e; \
echo "Running backend unit tests"; \
uv run --project api --dev pytest -p no:benchmark --timeout "$${PYTEST_TIMEOUT:-20}" -n auto \
api/tests/unit_tests \
+1 -1
View File
@@ -71,7 +71,7 @@ Dify is an open-source LLM app development platform. Its intuitive interface com
<br/>
The easiest way to start the Dify server is through [Docker Compose](docker/docker-compose.yaml). Before running Dify with the following commands, make sure that [Docker](https://docs.docker.com/get-docker/) and Docker Compose v2.24.0 or later are installed on your machine:
The easiest way to start the Dify server is through [Docker Compose](docker/docker-compose.yaml). Before running Dify with the following commands, make sure that [Docker](https://docs.docker.com/get-docker/) and [Docker Compose](https://docs.docker.com/compose/install/) are installed on your machine:
```bash
cd dify
+19 -1
View File
@@ -683,11 +683,28 @@ AGENT_BACKEND_RUN_TIMEOUT_SECONDS=1200
# KnowledgeFS (Dataset 2.0)
KNOWLEDGE_FS_ENABLED=false
# Production deployments require HTTPS; plain HTTP is limited to non-production or loopback.
KNOWLEDGE_FS_BASE_URL=
# Shared with KnowledgeFS; use at least 32 random characters.
KNOWLEDGE_FS_DIRECT_ORIGIN=
KNOWLEDGE_FS_LIFECYCLE_WORKER_ENABLED=false
KNOWLEDGE_FS_INTEGRATED_PROVISION_READY=false
KNOWLEDGE_FS_LEGACY_ACL_FREEZE_READY=false
KNOWLEDGE_FS_LIFECYCLE_POLL_INTERVAL_SECONDS=15
KNOWLEDGE_FS_LIFECYCLE_LEASE_SECONDS=60
KNOWLEDGE_FS_LIFECYCLE_BATCH_SIZE=25
# Legacy rollback-only HMAC; Capability v2 deployments leave this blank.
KNOWLEDGE_FS_JWT_SECRET=
KNOWLEDGE_FS_CAPABILITY_V2_ENABLED=false
KNOWLEDGE_FS_CAPABILITY_V2_SIGNING_KID=
KNOWLEDGE_FS_CAPABILITY_V2_PRIVATE_KEY_PEM=
KNOWLEDGE_FS_CAPABILITY_V2_PREVIOUS_PUBLIC_JWKS=
KNOWLEDGE_FS_CAPABILITY_V2_ISSUER=dify-control-plane
KNOWLEDGE_FS_CAPABILITY_V2_AUDIENCE=knowledge-fs
KNOWLEDGE_FS_CAPABILITY_V2_MAX_TTL_SECONDS=60
KNOWLEDGE_FS_SSE_READ_TIMEOUT_SECONDS=300
KNOWLEDGE_FS_TIMEOUT_SECONDS=10
KNOWLEDGE_FS_JWKS_CACHE_MAX_AGE_SECONDS=300
KNOWLEDGE_FS_PRODUCT_MAX_RESPONSE_BYTES=4194304
# Marketplace configuration
MARKETPLACE_ENABLED=true
@@ -729,6 +746,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
+2
View File
@@ -154,6 +154,7 @@ def initialize_extensions(app: DifyApp):
ext_forward_refs,
ext_hosting_provider,
ext_import_modules,
ext_knowledge_fs_observability,
ext_logging,
ext_login,
ext_logstore,
@@ -204,6 +205,7 @@ def initialize_extensions(app: DifyApp):
ext_enterprise_telemetry,
ext_request_logging,
ext_session_factory,
ext_knowledge_fs_observability,
ext_oauth_bearer,
]
for ext in extensions:
+2
View File
@@ -10,6 +10,7 @@ from .data_migration import (
import_migration_data,
migration_data_wizard,
)
from .knowledge_fs import knowledge_fs_control_space
from .plugin import (
backfill_plugin_auto_upgrade,
extract_plugins,
@@ -75,6 +76,7 @@ __all__ = [
"import_migration_data",
"install_plugins",
"install_rag_pipeline_plugins",
"knowledge_fs_control_space",
"legacy_model_types",
"migrate_annotation_vector_database",
"migrate_data_for_plugin",
+604
View File
@@ -0,0 +1,604 @@
"""Operator commands for the independent KnowledgeFS control-plane."""
from __future__ import annotations
import json
from collections.abc import Callable
from datetime import datetime
from functools import partial
from pathlib import Path
import click
from pydantic import BaseModel, ValidationError
from core.db.session_factory import session_factory
from services.knowledge_fs.cleanup import (
CleanupApprovalInput,
CleanupCompletionEvidenceInput,
CleanupReadinessEvidenceInput,
CleanupStartInput,
KnowledgeFSCleanupError,
KnowledgeFSCleanupService,
)
from services.knowledge_fs.control_space_commands import KnowledgeFSControlSpaceCommandService
from services.knowledge_fs.control_space_lifecycle import KnowledgeFSControlSpaceLifecycleError
from services.knowledge_fs.control_space_management import (
KnowledgeFSControlSpaceManagementService,
KnowledgeFSControlSpaceRegistration,
)
from services.knowledge_fs.cutover import (
CutoverSmokeResultsInput,
FinalDeltaInput,
KnowledgeFSCutoverError,
KnowledgeFSWorkspaceCutoverService,
LegacyDependencyInput,
QuarantineResolutionInput,
ShadowAuthorizationObservationInput,
ShadowCompletionInput,
WorkspaceInventoryInput,
)
from services.knowledge_fs.orphan_reconciler import KnowledgeFSOrphanReconciler
from services.knowledge_fs.remote_registry import get_knowledge_fs_lifecycle_remote
@click.group("knowledge-fs-control-space")
def knowledge_fs_control_space() -> None:
"""Inspect and repair Dify-owned KnowledgeFS control-space state."""
@knowledge_fs_control_space.command("dry-run")
@click.option("--tenant-id", default=None)
def dry_run(tenant_id: str | None) -> None:
report = _management_service().dry_run(tenant_id=tenant_id)
click.echo(json.dumps(report._asdict(), sort_keys=True))
@knowledge_fs_control_space.command("inventory")
@click.option("--input", "input_path", type=click.Path(path_type=Path, exists=True, dir_okay=False), required=True)
@click.option("--apply", is_flag=True, default=False, help="Create ledgers; omitted means read-only inventory.")
def inventory(input_path: Path, apply: bool) -> None:
"""Validate strict Workspace inventory JSONL and optionally create ledgers."""
service = _cutover_service()
for payload in _read_jsonl(input_path, WorkspaceInventoryInput):
report = _operator_call(partial(service.inventory, payload, apply=apply))
_echo_json(report._asdict())
@knowledge_fs_control_space.command("register")
@click.option("--tenant-id", required=True)
@click.option("--owner-account-id", required=True)
@click.option("--provisioning-key", required=True)
@click.option("--knowledge-space-id", required=True)
@click.option("--knowledge-space-revision", type=click.IntRange(min=0), required=True)
def register(
tenant_id: str,
owner_account_id: str,
provisioning_key: str,
knowledge_space_id: str,
knowledge_space_revision: int,
) -> None:
control_space, replayed = _management_service().register(
KnowledgeFSControlSpaceRegistration(
tenant_id,
owner_account_id,
provisioning_key,
knowledge_space_id,
knowledge_space_revision,
)
)
click.echo(json.dumps({"control_space_id": control_space.id, "replayed": replayed}, sort_keys=True))
@knowledge_fs_control_space.command("backfill")
@click.option("--input", "input_path", type=click.Path(path_type=Path, exists=True, dir_okay=False), required=True)
@click.option("--apply", is_flag=True, default=False, help="Persist registrations; omitted means dry-run.")
def backfill(input_path: Path, apply: bool) -> None:
"""Backfill strict Workspace inventory JSONL; dry-run unless --apply is explicit."""
service = _cutover_service()
for payload in _read_jsonl(input_path, WorkspaceInventoryInput):
report = _operator_call(partial(service.backfill, payload, apply=apply))
_echo_json(report._asdict())
@knowledge_fs_control_space.command("quarantine-resolve")
@click.option("--input", "input_path", type=click.Path(path_type=Path, exists=True, dir_okay=False), required=True)
@click.option("--apply", is_flag=True, default=False, help="Persist resolutions; omitted means dry-run.")
def quarantine_resolve(input_path: Path, apply: bool) -> None:
"""Resolve strict tenant-scoped quarantine JSONL with immutable operator evidence."""
service = _cutover_service()
for payload in _read_jsonl(input_path, QuarantineResolutionInput):
report = _operator_call(partial(service.resolve_quarantine, payload, apply=apply))
_echo_json(report._asdict())
@knowledge_fs_control_space.command("shadow-start")
@click.option("--tenant-id", required=True)
@click.option("--expected-cas-version", type=click.IntRange(min=0), required=True)
@click.option("--at", "started_at", default=None, help="Optional explicit timezone-aware shadow start.")
def shadow_start(tenant_id: str, expected_cas_version: int, started_at: str | None) -> None:
service = _cutover_service()
_operator_call(
lambda: service.begin_shadow(
tenant_id=tenant_id,
expected_cas_version=expected_cas_version,
started_at=_parse_timestamp(started_at) if started_at is not None else None,
)
)
_echo_json(service.status(tenant_id=tenant_id))
@knowledge_fs_control_space.command("shadow-report")
@click.option("--input", "input_path", type=click.Path(path_type=Path, exists=True, dir_okay=False), required=True)
@click.option("--apply", is_flag=True, default=False, help="Persist observations; omitted means dry-run.")
def shadow_report(input_path: Path, apply: bool) -> None:
observations = _read_jsonl(input_path, ShadowAuthorizationObservationInput)
report = _operator_call(lambda: _cutover_service().record_shadow_report(observations, apply=apply))
_echo_json(report._asdict())
@knowledge_fs_control_space.command("shadow-complete")
@click.option("--input", "input_path", type=click.Path(path_type=Path, exists=True, dir_okay=False), required=True)
@click.option("--apply", is_flag=True, default=False, help="Persist completion; omitted means dry-run.")
def shadow_complete(input_path: Path, apply: bool) -> None:
payload = _read_one_jsonl(input_path, ShadowCompletionInput)
report = _operator_call(lambda: _cutover_service().complete_shadow(payload, apply=apply))
_echo_json(report._asdict())
@knowledge_fs_control_space.command("issue-approve")
@click.option("--tenant-id", required=True)
@click.option("--issue-key", required=True)
@click.option("--account-id", required=True)
@click.option("--at", "approved_at", required=True)
def issue_approve(tenant_id: str, issue_key: str, account_id: str, approved_at: str) -> None:
_operator_call(
lambda: _cutover_service().approve_issue_fail_closed(
tenant_id=tenant_id,
issue_key=issue_key,
account_id=account_id,
approved_at=_parse_timestamp(approved_at),
)
)
_echo_json(_cutover_service().status(tenant_id=tenant_id))
@knowledge_fs_control_space.command("issue-resolve")
@click.option("--tenant-id", required=True)
@click.option("--issue-key", required=True)
@click.option("--account-id", required=True)
@click.option("--at", "resolved_at", required=True)
def issue_resolve(tenant_id: str, issue_key: str, account_id: str, resolved_at: str) -> None:
_operator_call(
lambda: _cutover_service().resolve_issue(
tenant_id=tenant_id,
issue_key=issue_key,
account_id=account_id,
resolved_at=_parse_timestamp(resolved_at),
)
)
_echo_json(_cutover_service().status(tenant_id=tenant_id))
@knowledge_fs_control_space.command("shadow-approve")
@click.option("--tenant-id", required=True)
@click.option("--diff-key", required=True)
@click.option("--account-id", required=True)
@click.option("--at", "approved_at", required=True)
def shadow_approve(tenant_id: str, diff_key: str, account_id: str, approved_at: str) -> None:
_operator_call(
lambda: _cutover_service().approve_shadow_diff(
tenant_id=tenant_id,
diff_key=diff_key,
account_id=account_id,
approved_at=_parse_timestamp(approved_at),
)
)
_echo_json(_cutover_service().status(tenant_id=tenant_id))
@knowledge_fs_control_space.command("shadow-resolve")
@click.option("--tenant-id", required=True)
@click.option("--diff-key", required=True)
@click.option("--account-id", required=True)
@click.option("--at", "resolved_at", required=True)
def shadow_resolve(tenant_id: str, diff_key: str, account_id: str, resolved_at: str) -> None:
_operator_call(
lambda: _cutover_service().resolve_shadow_diff(
tenant_id=tenant_id,
diff_key=diff_key,
account_id=account_id,
resolved_at=_parse_timestamp(resolved_at),
)
)
_echo_json(_cutover_service().status(tenant_id=tenant_id))
@knowledge_fs_control_space.command("legacy-dashboard")
@click.option("--tenant-id", required=True)
@click.option("--input", "input_path", type=click.Path(path_type=Path, exists=True, dir_okay=False), required=True)
@click.option("--checked-at", required=True)
@click.option("--expected-cas-version", type=click.IntRange(min=0), default=None)
@click.option("--apply", is_flag=True, default=False, help="Persist gate evidence; omitted means read-only report.")
def legacy_dashboard(
tenant_id: str,
input_path: Path,
checked_at: str,
expected_cas_version: int | None,
apply: bool,
) -> None:
dependencies = _read_jsonl(input_path, LegacyDependencyInput, allow_empty=True)
report = _operator_call(
lambda: _cutover_service().legacy_dependency_dashboard(
tenant_id=tenant_id,
dependencies=dependencies,
expected_cas_version=expected_cas_version,
checked_at=_parse_timestamp(checked_at),
apply=apply,
)
)
_echo_json(report._asdict())
@knowledge_fs_control_space.command("legacy-check")
@click.option("--tenant-id", required=True)
def legacy_check(tenant_id: str) -> None:
status_report = _operator_call(lambda: _cutover_service().status(tenant_id=tenant_id))
passed = (
bool(status_report["legacy_dependency_ready"])
and status_report["open_issues"] == 0
and status_report["unresolved_cutover_quarantine"] == 0
)
_echo_json({"tenant_id": tenant_id, "passed": passed, "status": status_report})
if not passed:
raise click.exceptions.Exit(1)
@knowledge_fs_control_space.command("freeze")
@click.option("--tenant-id", required=True)
@click.option("--expected-cas-version", type=click.IntRange(min=0), required=True)
@click.option("--at", "freeze_at", required=True)
def freeze(tenant_id: str, expected_cas_version: int, freeze_at: str) -> None:
service = _cutover_service()
_operator_call(
lambda: service.freeze(
tenant_id=tenant_id,
expected_cas_version=expected_cas_version,
freeze_at=_parse_timestamp(freeze_at),
)
)
_echo_json(service.status(tenant_id=tenant_id))
@knowledge_fs_control_space.command("final-delta")
@click.option("--input", "input_path", type=click.Path(path_type=Path, exists=True, dir_okay=False), required=True)
def final_delta(input_path: Path) -> None:
payload = _read_one_jsonl(input_path, FinalDeltaInput)
service = _cutover_service()
_operator_call(lambda: service.apply_final_delta(payload))
_echo_json(service.status(tenant_id=str(payload.tenant_id)))
@knowledge_fs_control_space.command("cutover")
@click.option("--tenant-id", required=True)
@click.option("--expected-cas-version", type=click.IntRange(min=0), required=True)
@click.option("--at", "cutover_at", required=True)
@click.option("--rollback-cutoff-at", required=True)
def cutover(tenant_id: str, expected_cas_version: int, cutover_at: str, rollback_cutoff_at: str) -> None:
service = _cutover_service()
_operator_call(
lambda: service.cutover(
tenant_id=tenant_id,
expected_cas_version=expected_cas_version,
cutover_at=_parse_timestamp(cutover_at),
rollback_cutoff_at=_parse_timestamp(rollback_cutoff_at),
)
)
_echo_json(service.status(tenant_id=tenant_id))
@knowledge_fs_control_space.command("smoke")
@click.option("--tenant-id", required=True)
@click.option("--expected-cas-version", type=click.IntRange(min=0), required=True)
@click.option("--input", "input_path", type=click.Path(path_type=Path, exists=True, dir_okay=False), required=True)
def smoke(tenant_id: str, expected_cas_version: int, input_path: Path) -> None:
results = _read_one_jsonl(input_path, CutoverSmokeResultsInput)
service = _cutover_service()
_operator_call(
lambda: service.record_smoke_results(
tenant_id=tenant_id,
expected_cas_version=expected_cas_version,
results=results,
)
)
_echo_json(service.status(tenant_id=tenant_id))
@knowledge_fs_control_space.command("observe")
@click.option("--tenant-id", required=True)
@click.option("--expected-cas-version", type=click.IntRange(min=0), required=True)
@click.option("--started-at", default=None)
@click.option("--window-ends-at", default=None)
@click.option("--maximum-task-expires-at", default=None)
@click.option("--observed-at", default=None)
def observe(
tenant_id: str,
expected_cas_version: int,
started_at: str | None,
window_ends_at: str | None,
maximum_task_expires_at: str | None,
observed_at: str | None,
) -> None:
service = _cutover_service()
if observed_at is not None:
if any(value is not None for value in (started_at, window_ends_at, maximum_task_expires_at)):
raise click.UsageError("--observed-at cannot be combined with observation start options")
_operator_call(
lambda: service.complete_observation(
tenant_id=tenant_id,
expected_cas_version=expected_cas_version,
observed_at=_parse_timestamp(observed_at),
)
)
else:
if started_at is None or window_ends_at is None or maximum_task_expires_at is None:
raise click.UsageError(
"observation start requires --started-at, --window-ends-at, and --maximum-task-expires-at"
)
_operator_call(
lambda: service.begin_observation(
tenant_id=tenant_id,
expected_cas_version=expected_cas_version,
started_at=_parse_timestamp(started_at),
window_ends_at=_parse_timestamp(window_ends_at),
maximum_task_expires_at=_parse_timestamp(maximum_task_expires_at),
)
)
_echo_json(service.status(tenant_id=tenant_id))
@knowledge_fs_control_space.command("rollback")
@click.option("--tenant-id", required=True)
@click.option("--expected-cas-version", type=click.IntRange(min=0), required=True)
@click.option("--at", "rolled_back_at", required=True)
def rollback(tenant_id: str, expected_cas_version: int, rolled_back_at: str) -> None:
service = _cutover_service()
_operator_call(
lambda: service.rollback(
tenant_id=tenant_id,
expected_cas_version=expected_cas_version,
rolled_back_at=_parse_timestamp(rolled_back_at),
)
)
_echo_json(service.status(tenant_id=tenant_id))
@knowledge_fs_control_space.command("status")
@click.option("--tenant-id", required=True)
def status(tenant_id: str) -> None:
_echo_json(_operator_call(lambda: _cutover_service().status(tenant_id=tenant_id)))
@knowledge_fs_control_space.command("cleanup-request")
@click.option("--input", "input_path", type=click.Path(path_type=Path, exists=True, dir_okay=False), required=True)
@click.option("--apply", is_flag=True, default=False, help="Persist readiness evidence; omitted means dry-run.")
def cleanup_request(input_path: Path, apply: bool) -> None:
payload = _read_one_jsonl(input_path, CleanupReadinessEvidenceInput)
report = _cleanup_call(lambda: _cleanup_service().request_cleanup(payload, apply=apply))
_echo_json(report._asdict())
@knowledge_fs_control_space.command("cleanup-approve")
@click.option("--input", "input_path", type=click.Path(path_type=Path, exists=True, dir_okay=False), required=True)
@click.option("--apply", is_flag=True, default=False, help="Persist four-eyes approval; omitted means dry-run.")
def cleanup_approve(input_path: Path, apply: bool) -> None:
payload = _read_one_jsonl(input_path, CleanupApprovalInput)
report = _cleanup_call(lambda: _cleanup_service().approve_cleanup(payload, apply=apply))
_echo_json(report._asdict())
@knowledge_fs_control_space.command("cleanup-start")
@click.option("--input", "input_path", type=click.Path(path_type=Path, exists=True, dir_okay=False), required=True)
@click.option("--apply", is_flag=True, default=False, help="Persist the irreversible fence; never runs deletion.")
@click.option(
"--acknowledge-irreversible",
is_flag=True,
default=False,
help="Required with --apply; confirms rollback will be permanently closed.",
)
def cleanup_start(input_path: Path, apply: bool, acknowledge_irreversible: bool) -> None:
payload = _read_one_jsonl(input_path, CleanupStartInput)
if apply and not acknowledge_irreversible:
raise click.UsageError("--apply requires --acknowledge-irreversible")
report = _cleanup_call(lambda: _cleanup_service().start_cleanup(payload, apply=apply))
_echo_json(report._asdict())
@knowledge_fs_control_space.command("cleanup-complete")
@click.option("--input", "input_path", type=click.Path(path_type=Path, exists=True, dir_okay=False), required=True)
@click.option("--apply", is_flag=True, default=False, help="Persist externally verified cleanup completion.")
@click.option(
"--acknowledge-executed",
is_flag=True,
default=False,
help="Required with --apply; confirms the reviewed destructive bundle already executed.",
)
def cleanup_complete(input_path: Path, apply: bool, acknowledge_executed: bool) -> None:
payload = _read_one_jsonl(input_path, CleanupCompletionEvidenceInput)
if apply and not acknowledge_executed:
raise click.UsageError("--apply requires --acknowledge-executed")
report = _cleanup_call(lambda: _cleanup_service().complete_cleanup(payload, apply=apply))
_echo_json(report._asdict())
@knowledge_fs_control_space.command("cleanup-status")
@click.option("--tenant-id", required=True)
@click.option("--request-id", required=True)
def cleanup_status(tenant_id: str, request_id: str) -> None:
_echo_json(_cleanup_call(lambda: _cleanup_service().status(tenant_id=tenant_id, request_id=request_id)))
@knowledge_fs_control_space.command("repair")
@click.option("--tenant-id", required=True)
@click.option("--control-space-id", required=True)
@click.option("--expected-resource-version", type=click.IntRange(min=0), required=True)
@click.option("--knowledge-space-id", required=True)
@click.option("--knowledge-space-revision", type=click.IntRange(min=0), required=True)
def repair(
tenant_id: str,
control_space_id: str,
expected_resource_version: int,
knowledge_space_id: str,
knowledge_space_revision: int,
) -> None:
control_space = _management_service().repair_registration(
tenant_id=tenant_id,
control_space_id=control_space_id,
expected_resource_version=expected_resource_version,
knowledge_space_id=knowledge_space_id,
knowledge_space_revision=knowledge_space_revision,
)
click.echo(json.dumps({"control_space_id": control_space.id, "state": control_space.state.value}, sort_keys=True))
@knowledge_fs_control_space.command("orphan-report")
@click.option("--limit", type=click.IntRange(min=1, max=10_000), default=500, show_default=True)
def orphan_report(limit: int) -> None:
report = KnowledgeFSOrphanReconciler(
session_factory.get_session_maker(),
get_knowledge_fs_lifecycle_remote(),
).reconcile(limit=limit, apply_repairs=False)
click.echo(json.dumps(report._asdict(), sort_keys=True))
@knowledge_fs_control_space.command("workspace-delete-request")
@click.option("--tenant-id", required=True)
@click.option("--apply", is_flag=True, default=False, help="Persist durable deletion intents; omitted means dry-run.")
def workspace_delete_request(tenant_id: str, apply: bool) -> None:
"""Route every KnowledgeFS Space through the canonical lifecycle deletion path."""
if not apply:
report = _management_service().dry_run(tenant_id=tenant_id)
_echo_json(
{
"apply": False,
"by_state": report.by_state,
"tenant_id": tenant_id,
"total": report.total,
}
)
return
results = _lifecycle_call(lambda: _lifecycle_service().request_workspace_cleanup(tenant_id=tenant_id))
_echo_json(
{
"apply": True,
"control_space_ids": [result.control_space.id for result in results],
"operation_ids": [result.outbox.operation_id for result in results if result.outbox is not None],
"tenant_id": tenant_id,
}
)
@knowledge_fs_control_space.command("workspace-delete-finalize")
@click.option("--tenant-id", required=True)
@click.option("--apply", is_flag=True, default=False, help="Purge terminal local control-plane rows.")
@click.option(
"--acknowledge-control-plane-purge",
is_flag=True,
default=False,
help="Required with --apply after every remote Space has reached deleted.",
)
def workspace_delete_finalize(tenant_id: str, apply: bool, acknowledge_control_plane_purge: bool) -> None:
"""Release the Workspace FK only after all remote deletions are terminal."""
service = _lifecycle_service()
if not apply:
_lifecycle_call(lambda: service.assert_workspace_deletion_allowed(tenant_id=tenant_id))
_echo_json({"apply": False, "ready": True, "tenant_id": tenant_id})
return
if not acknowledge_control_plane_purge:
raise click.UsageError("--apply requires --acknowledge-control-plane-purge")
deleted = _lifecycle_call(lambda: service.finalize_workspace_deletion(tenant_id=tenant_id))
_echo_json({"apply": True, "purged_control_spaces": deleted, "tenant_id": tenant_id})
def _read_jsonl[InputT: BaseModel](
input_path: Path, input_type: type[InputT], *, allow_empty: bool = False
) -> tuple[InputT, ...]:
records: list[InputT] = []
for line_number, line in enumerate(input_path.read_text(encoding="utf-8").splitlines(), start=1):
if not line.strip():
continue
try:
records.append(input_type.model_validate_json(line))
except ValidationError as exc:
raise click.ClickException(f"invalid strict JSONL at line {line_number}: {exc}") from exc
if not records and not allow_empty:
raise click.ClickException("strict JSONL input must contain at least one record")
return tuple(records)
def _read_one_jsonl[InputT: BaseModel](input_path: Path, input_type: type[InputT]) -> InputT:
records = _read_jsonl(input_path, input_type)
if len(records) != 1:
raise click.ClickException("this command requires exactly one JSONL record")
return records[0]
def _parse_timestamp(value: str) -> datetime:
try:
parsed = datetime.fromisoformat(value)
except ValueError as exc:
raise click.ClickException(f"invalid ISO-8601 timestamp: {value}") from exc
if parsed.tzinfo is None:
raise click.ClickException("operator timestamps must include an explicit timezone")
return parsed
def _operator_call[ResultT](operation: Callable[[], ResultT]) -> ResultT:
try:
return operation()
except KnowledgeFSCutoverError as exc:
raise click.ClickException(str(exc)) from exc
def _cleanup_call[ResultT](operation: Callable[[], ResultT]) -> ResultT:
try:
return operation()
except KnowledgeFSCleanupError as exc:
raise click.ClickException(str(exc)) from exc
def _lifecycle_call[ResultT](operation: Callable[[], ResultT]) -> ResultT:
try:
return operation()
except KnowledgeFSControlSpaceLifecycleError as exc:
raise click.ClickException(str(exc)) from exc
def _echo_json(payload: object) -> None:
click.echo(json.dumps(payload, default=str, sort_keys=True))
def _management_service() -> KnowledgeFSControlSpaceManagementService:
return KnowledgeFSControlSpaceManagementService(session_factory.get_session_maker())
def _cutover_service() -> KnowledgeFSWorkspaceCutoverService:
return KnowledgeFSWorkspaceCutoverService(
session_factory.get_session_maker(),
remote_factory=get_knowledge_fs_lifecycle_remote,
)
def _cleanup_service() -> KnowledgeFSCleanupService:
return KnowledgeFSCleanupService(session_factory.get_session_maker())
def _lifecycle_service() -> KnowledgeFSControlSpaceCommandService:
return KnowledgeFSControlSpaceCommandService(session_factory.get_session_maker())
__all__ = ["knowledge_fs_control_space"]
-6
View File
@@ -14,12 +14,6 @@ class EnterpriseFeatureConfig(BaseSettings):
default=False,
)
WEBAPP_PUBLIC_ACCESS_ENABLED: bool = Field(
description="Whether admins are allowed to set a webapp's access mode to public (anyone with the link, "
"no auth). Disable in security-sensitive on-prem deployments.",
default=True,
)
CAN_REPLACE_LOGO: bool = Field(
description="Allow customization of the enterprise logo.",
default=False,
+90 -20
View File
@@ -1,30 +1,68 @@
"""Configuration for the optional KnowledgeFS Console bridge."""
"""Configuration for the optional KnowledgeFS control-plane integration."""
from ipaddress import ip_address
from urllib.parse import urlsplit
from pydantic import Field, PositiveFloat, SecretStr, field_validator, model_validator
from pydantic import Field, PositiveFloat, PositiveInt, SecretStr, field_validator, model_validator
from pydantic_settings import BaseSettings
class KnowledgeFSConfig(BaseSettings):
"""Server-only settings for the KnowledgeFS production connection."""
"""Server-only KnowledgeFS connection and rollout settings."""
KNOWLEDGE_FS_ENABLED: bool = Field(
default=False,
description="Enable the private KnowledgeFS Console bridge.",
)
KNOWLEDGE_FS_BASE_URL: str | None = Field(default=None, description="KnowledgeFS gateway base URL.")
KNOWLEDGE_FS_JWT_SECRET: SecretStr | None = Field(
default=None,
min_length=32,
description="Shared secret used to sign short-lived KnowledgeFS service JWTs.",
KNOWLEDGE_FS_LIFECYCLE_WORKER_ENABLED: bool = Field(
default=False,
description="Enable delivery of durable KnowledgeFS lifecycle commands after every rollout gate is ready.",
)
KNOWLEDGE_FS_SSE_READ_TIMEOUT_SECONDS: PositiveFloat = Field(default=300.0, le=3600.0, allow_inf_nan=False)
KNOWLEDGE_FS_INTEGRATED_PROVISION_READY: bool = Field(
default=False,
description="Confirm that the Capability-v2 integrated provision route is deployed and verified.",
)
KNOWLEDGE_FS_LEGACY_ACL_FREEZE_READY: bool = Field(
default=False,
description="Confirm that legacy KFS ACL mutation is frozen for integrated mode.",
)
KNOWLEDGE_FS_LIFECYCLE_POLL_INTERVAL_SECONDS: PositiveInt = Field(default=15, le=300)
KNOWLEDGE_FS_LIFECYCLE_LEASE_SECONDS: PositiveInt = Field(default=60, le=600)
KNOWLEDGE_FS_LIFECYCLE_BATCH_SIZE: PositiveInt = Field(default=25, le=1_000)
KNOWLEDGE_FS_BASE_URL: str | None = Field(default=None, description="KnowledgeFS gateway base URL.")
KNOWLEDGE_FS_DIRECT_ORIGIN: str | None = Field(
default=None,
description="Public KnowledgeFS origin returned with direct upload capabilities.",
)
KNOWLEDGE_FS_CAPABILITY_V2_ENABLED: bool = Field(
default=False,
description="Prepare resource-scoped Capability v2 issuance; disabled until rollout approval.",
)
KNOWLEDGE_FS_CAPABILITY_V2_SIGNING_KID: str | None = Field(
default=None,
description="Identifier for the current asymmetric Capability v2 signing key.",
)
KNOWLEDGE_FS_CAPABILITY_V2_PRIVATE_KEY_PEM: SecretStr | None = Field(
default=None,
description="Server-only PEM for the current Capability v2 RSA signing key.",
)
KNOWLEDGE_FS_CAPABILITY_V2_PREVIOUS_PUBLIC_JWKS: str | None = Field(
default=None,
description="Optional public-only JWKS JSON retained during key rotation overlap.",
)
KNOWLEDGE_FS_CAPABILITY_V2_ISSUER: str = Field(default="dify-control-plane", min_length=1)
KNOWLEDGE_FS_CAPABILITY_V2_AUDIENCE: str = Field(default="knowledge-fs", min_length=1)
KNOWLEDGE_FS_CAPABILITY_V2_MAX_TTL_SECONDS: PositiveInt = Field(default=60, le=60)
KNOWLEDGE_FS_JWKS_CACHE_MAX_AGE_SECONDS: PositiveInt = Field(default=300, le=86_400)
KNOWLEDGE_FS_PRODUCT_MAX_RESPONSE_BYTES: PositiveInt = Field(default=4 * 1024 * 1024, le=16 * 1024 * 1024)
KNOWLEDGE_FS_TIMEOUT_SECONDS: PositiveFloat = Field(default=10.0, le=60.0, allow_inf_nan=False)
@field_validator(
"KNOWLEDGE_FS_BASE_URL",
"KNOWLEDGE_FS_JWT_SECRET",
"KNOWLEDGE_FS_DIRECT_ORIGIN",
"KNOWLEDGE_FS_CAPABILITY_V2_SIGNING_KID",
"KNOWLEDGE_FS_CAPABILITY_V2_PRIVATE_KEY_PEM",
"KNOWLEDGE_FS_CAPABILITY_V2_PREVIOUS_PUBLIC_JWKS",
mode="before",
)
@classmethod
@@ -40,25 +78,57 @@ class KnowledgeFSConfig(BaseSettings):
@field_validator("KNOWLEDGE_FS_BASE_URL")
@classmethod
def validate_base_url(cls, value: str | None) -> str | None:
return cls._validate_origin(value, name="KNOWLEDGE_FS_BASE_URL")
@field_validator("KNOWLEDGE_FS_DIRECT_ORIGIN")
@classmethod
def validate_direct_origin(cls, value: str | None) -> str | None:
return cls._validate_origin(value, name="KNOWLEDGE_FS_DIRECT_ORIGIN")
@classmethod
def _validate_origin(cls, value: str | None, *, name: str) -> str | None:
if value is None:
return None
parsed = urlsplit(value)
if parsed.scheme not in {"http", "https"} or not parsed.netloc:
raise ValueError("KNOWLEDGE_FS_BASE_URL must be an absolute HTTP(S) URL")
raise ValueError(f"{name} must be an absolute HTTP(S) URL")
try:
_ = parsed.port
except ValueError as exc:
raise ValueError("KNOWLEDGE_FS_BASE_URL must include a valid port") from exc
if parsed.username or parsed.password or parsed.query or parsed.fragment:
raise ValueError("KNOWLEDGE_FS_BASE_URL must not include credentials, query, or fragment")
raise ValueError(f"{name} must include a valid port") from exc
if parsed.username or parsed.password or parsed.query or parsed.fragment or parsed.path not in {"", "/"}:
raise ValueError(f"{name} must be an origin without credentials, path, query, or fragment")
return value.rstrip("/")
@model_validator(mode="after")
def validate_enabled_connection(self) -> "KnowledgeFSConfig":
if not self.KNOWLEDGE_FS_ENABLED:
return self
if bool(self.KNOWLEDGE_FS_BASE_URL) != bool(self.KNOWLEDGE_FS_JWT_SECRET):
raise ValueError("KNOWLEDGE_FS_BASE_URL and KNOWLEDGE_FS_JWT_SECRET must be configured together")
if not self.KNOWLEDGE_FS_BASE_URL:
raise ValueError("KnowledgeFS connection settings are required when the integration is enabled")
if str(getattr(self, "DEPLOY_ENV", "")).strip().upper() == "PRODUCTION":
for name, value in (
("KNOWLEDGE_FS_BASE_URL", self.KNOWLEDGE_FS_BASE_URL),
("KNOWLEDGE_FS_DIRECT_ORIGIN", self.KNOWLEDGE_FS_DIRECT_ORIGIN),
):
if value and not self._is_secure_or_loopback_origin(value):
raise ValueError(f"{name} must use HTTPS in production unless it targets loopback")
if self.KNOWLEDGE_FS_ENABLED:
if not self.KNOWLEDGE_FS_BASE_URL:
raise ValueError("KnowledgeFS base URL is required when the integration is enabled")
if not self.KNOWLEDGE_FS_CAPABILITY_V2_ENABLED:
raise ValueError("KnowledgeFS product routes require Capability v2 when enabled")
if self.KNOWLEDGE_FS_CAPABILITY_V2_ENABLED and not (
self.KNOWLEDGE_FS_CAPABILITY_V2_SIGNING_KID and self.KNOWLEDGE_FS_CAPABILITY_V2_PRIVATE_KEY_PEM
):
raise ValueError("Capability v2 signing kid and private key are required when issuance is enabled")
return self
@staticmethod
def _is_secure_or_loopback_origin(value: str) -> bool:
parsed = urlsplit(value)
if parsed.scheme == "https":
return True
hostname = (parsed.hostname or "").rstrip(".").lower()
if hostname == "localhost":
return True
try:
return ip_address(hostname).is_loopback
except ValueError:
return False
+1 -69
View File
@@ -1,4 +1,4 @@
from datetime import datetime, timedelta
from datetime import timedelta
from enum import StrEnum
from typing import Literal
@@ -11,7 +11,6 @@ from pydantic import (
PositiveFloat,
PositiveInt,
computed_field,
field_validator,
)
from pydantic_settings import BaseSettings
@@ -281,27 +280,6 @@ class PluginConfig(BaseSettings):
default="",
)
@field_validator("PLUGIN_REMOTE_INSTALL_PORT", mode="before")
@classmethod
def _reject_host_port_shaped_plugin_remote_install_port(cls, v):
"""Reject ``host:port``-shaped values with an actionable hint.
``EXPOSE_PLUGIN_DEBUGGING_PORT`` is overloaded: it feeds both the
plugin_daemon ``ports:`` mapping (where ``127.0.0.1:5003`` is valid
compose syntax) and this integer app setting advertised in the console.
Without this guard a loopback bind spec crashloops the api container
with an opaque ``int_parsing`` traceback. See issue #39323.
"""
if isinstance(v, str) and ":" in v.strip():
raise ValueError(
"PLUGIN_REMOTE_INSTALL_PORT must be a bare port number, got "
f"{v!r}. A 'host:port' value usually means "
"EXPOSE_PLUGIN_DEBUGGING_PORT was set to a compose publish spec "
"like '127.0.0.1:5003'; bind loopback via a "
"docker-compose.override.yaml instead of overloading this var."
)
return v
@property
def NEW_USER_DEFAULT_PLUGIN_ID_LIST(self) -> list[str]:
return [item.strip() for item in self.NEW_USER_DEFAULT_PLUGIN_IDS.split(",") if item.strip()]
@@ -816,41 +794,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
@@ -1195,16 +1138,6 @@ class HomepageConfig(BaseSettings):
default=True,
)
ENABLE_STEP_BY_STEP_TOUR: bool = Field(
description="Enable account-level Step-by-step Tour eligibility checks",
default=False,
)
STEP_BY_STEP_TOUR_ROLLOUT_STARTED_AT: datetime | None = Field(
description="UTC timestamp after which newly initialized accounts are eligible for Step-by-step Tour",
default=None,
)
class RagEtlConfig(BaseSettings):
"""
@@ -1629,7 +1562,6 @@ class FeatureConfig(
TenantIsolatedTaskQueueConfig,
ToolConfig,
UpdateConfig,
CommunityTelemetryConfig,
WorkflowConfig,
WorkflowNodeExecutionConfig,
WorkspaceConfig,
+2 -4
View File
@@ -38,9 +38,7 @@ from . import (
feature,
human_input_form,
init_validate,
knowledge_fs_proxy,
notification,
onboarding,
ping,
setup,
spec,
@@ -127,6 +125,7 @@ from .explore import (
saved_message,
trial,
)
from .knowledge_fs import resources as knowledge_fs_resources
from .snippets import snippet_workflow, snippet_workflow_draft_variable
from .socketio import workflow as socketio_workflow
@@ -197,7 +196,7 @@ __all__ = [
"human_input_form",
"init_validate",
"installed_app",
"knowledge_fs_proxy",
"knowledge_fs_resources",
"load_balancing_config",
"login",
"mcp_server",
@@ -210,7 +209,6 @@ __all__ = [
"notification",
"oauth",
"oauth_server",
"onboarding",
"ops_trace",
"parameter",
"ping",
+1 -24
View File
@@ -62,7 +62,7 @@ from libs.datetime_utils import parse_time_range
from libs.helper import dump_response
from libs.login import login_required
from models import Account
from models.agent import Agent, AgentConfigDraftType, AgentStatus
from models.agent import Agent, AgentStatus
from models.agent_config_entities import AgentSoulConfig
from models.enums import ApiTokenType
from models.model import ApiToken, App, IconType
@@ -266,13 +266,6 @@ class AgentDebugConversationRefreshResponse(BaseModel):
debug_conversation_message_count: int = 0
class AgentDebugConversationRefreshPayload(BaseModel):
draft_type: AgentConfigDraftType = Field(
default=AgentConfigDraftType.DEBUG_BUILD,
description="Agent draft surface whose conversation should be refreshed",
)
class AgentPublishPayload(BaseModel):
version_note: str | None = Field(default=None, description="Optional note for this published Agent version")
@@ -316,7 +309,6 @@ register_schema_models(
AgentAppCopyPayload,
AgentPublishPayload,
AgentBuildDraftCheckoutPayload,
AgentDebugConversationRefreshPayload,
ComposerSavePayload,
AgentApiStatusPayload,
AgentInviteOptionsQuery,
@@ -400,7 +392,6 @@ def _serialize_agent_app_detail(
tenant_id=app_model.tenant_id,
agent_id=agent.id,
account_id=current_user.id,
draft_type=AgentConfigDraftType.DEBUG_BUILD,
commit=False,
)
message_count = roster_service.count_agent_app_debug_conversation_messages(
@@ -448,7 +439,6 @@ def _serialize_agent_app_pagination(session: Session, app_pagination, *, tenant_
tenant_id=tenant_id,
agents=list(agents_by_app_id.values()),
account_id=current_user.id,
draft_type=AgentConfigDraftType.DEBUG_BUILD,
)
payload = AgentAppPagination.model_validate(
app_pagination,
@@ -665,16 +655,6 @@ class AgentAppApi(Resource):
@console_ns.route("/agent/<uuid:agent_id>/debug-conversation/refresh")
class AgentDebugConversationRefreshApi(Resource):
@console_ns.expect(console_ns.models[AgentDebugConversationRefreshPayload.__name__])
@console_ns.doc(
params={
"payload": {
"in": "body",
"required": False,
"schema": {"$ref": f"#/components/schemas/{AgentDebugConversationRefreshPayload.__name__}"},
}
}
)
@console_ns.response(
200,
"Agent debug conversation refreshed",
@@ -689,12 +669,10 @@ class AgentDebugConversationRefreshApi(Resource):
@with_current_tenant_id
@with_session
def post(self, session: Session, tenant_id: str, current_user: Account, agent_id: UUID):
args = AgentDebugConversationRefreshPayload.model_validate(request.get_json(silent=True) or {})
debug_conversation_id = _agent_roster_service(session).refresh_agent_app_debug_conversation_id(
tenant_id=tenant_id,
agent_id=str(agent_id),
account_id=current_user.id,
draft_type=args.draft_type,
)
return AgentDebugConversationRefreshResponse(
debug_conversation_id=debug_conversation_id,
@@ -751,7 +729,6 @@ class AgentBuildDraftCheckoutApi(Resource):
@console_ns.route("/agent/<uuid:agent_id>/build-draft")
class AgentBuildDraftApi(Resource):
@console_ns.response(200, "Agent build draft", console_ns.models[AgentBuildDraftResponse.__name__])
@console_ns.response(404, "Agent build draft not found")
@setup_required
@login_required
@account_initialization_required
+5 -9
View File
@@ -246,7 +246,7 @@ class ModelConfigPartial(ResponseModel):
return to_timestamp(value)
class AppModelConfigResponse(ResponseModel):
class ModelConfig(ResponseModel):
opening_statement: str | None = None
suggested_questions: Any | None = Field(
default=None, validation_alias=AliasChoices("suggested_questions_list", "suggested_questions")
@@ -419,7 +419,7 @@ class AppDetail(AppResponseModel):
icon_background: str | None = None
enable_site: bool
enable_api: bool
model_config_: AppModelConfigResponse | None = Field(
model_config_: ModelConfig | None = Field(
default=None,
validation_alias=AliasChoices("app_model_config", "model_config"),
alias="model_config",
@@ -525,13 +525,7 @@ def _enrich_app_list_items(session: Session, *, apps: Sequence[App], tenant_id:
register_enum_models(console_ns, RetrievalMethod, WorkflowExecutionStatus, DatasetPermissionEnum)
register_response_schema_models(
console_ns,
RedirectUrlResponse,
SimpleResultResponse,
AppImportResponse,
AppTraceResponse,
AppModelConfigResponse,
AppDetail,
console_ns, RedirectUrlResponse, SimpleResultResponse, AppImportResponse, AppTraceResponse
)
register_schema_models(
@@ -550,8 +544,10 @@ register_schema_models(
Tag,
WorkflowPartial,
ModelConfigPartial,
ModelConfig,
AppDetailSiteResponse,
DeletedTool,
AppDetail,
AppExportResponse,
Segmentation,
PreProcessingRule,
+1 -14
View File
@@ -49,7 +49,6 @@ from libs import helper
from libs.helper import uuid_value
from libs.login import login_required
from models import Account
from models.agent import AgentConfigDraftType
from models.model import App, AppMode
from services.agent.errors import AgentNotFoundError
from services.agent.roster_service import AgentRosterService
@@ -344,23 +343,14 @@ class AgentChatMessageStopApi(Resource):
def _resolve_current_user_agent_debug_conversation_id(
*,
session: Session,
current_tenant_id: str,
current_user: Account,
app_model: App,
agent_id: str | None,
draft_type: AgentConfigDraftType,
*, session: Session, current_tenant_id: str, current_user: Account, app_model: App, agent_id: str | None
) -> str:
"""Resolve the current editor's conversation without crossing draft surfaces."""
roster_service = AgentRosterService(session)
if agent_id:
return roster_service.get_or_create_agent_app_debug_conversation_id(
tenant_id=current_tenant_id,
agent_id=agent_id,
account_id=current_user.id,
draft_type=draft_type,
)
agent = roster_service.get_app_backing_agent(tenant_id=current_tenant_id, app_id=str(app_model.id))
@@ -370,7 +360,6 @@ def _resolve_current_user_agent_debug_conversation_id(
tenant_id=current_tenant_id,
agent_id=agent.id,
account_id=current_user.id,
draft_type=draft_type,
)
@@ -393,7 +382,6 @@ def _create_chat_message(
current_user=current_user,
app_model=app_model,
agent_id=agent_id,
draft_type=AgentConfigDraftType(args_model.draft_type),
)
if args_model.conversation_id and args_model.conversation_id != debug_conversation_id:
raise NotFound("Conversation Not Exists.")
@@ -430,7 +418,6 @@ def _create_build_chat_finalization_message(
current_user=current_user,
app_model=app_model,
agent_id=agent_id,
draft_type=AgentConfigDraftType.DEBUG_BUILD,
)
args: dict[str, Any] = {
"query": _BUILD_CHAT_FINALIZATION_QUERY,
+13 -16
View File
@@ -359,12 +359,6 @@ class WorkflowPublishResponse(ResponseModel):
created_at: int
class SyncDraftWorkflowResponse(ResponseModel):
result: str
hash: str
updated_at: int
class WorkflowRestoreResponse(ResponseModel):
result: str
hash: str
@@ -447,7 +441,6 @@ register_response_schema_models(
WorkflowOnlineUsersByApp,
WorkflowOnlineUsersResponse,
WorkflowPublishResponse,
SyncDraftWorkflowResponse,
WorkflowRestoreResponse,
DefaultBlockConfigsResponse,
DefaultBlockConfigResponse,
@@ -563,7 +556,14 @@ class DraftWorkflowApi(Resource):
@console_ns.response(
200,
"Draft workflow synced successfully",
console_ns.models[SyncDraftWorkflowResponse.__name__],
console_ns.model(
"SyncDraftWorkflowResponse",
{
"result": fields.String,
"hash": fields.String,
"updated_at": fields.String,
},
),
)
@console_ns.response(400, "Invalid workflow configuration")
@console_ns.response(403, "Permission denied")
@@ -618,14 +618,11 @@ class DraftWorkflowApi(Resource):
except VariableError as e:
raise InvalidArgumentError(description=str(e))
return dump_response(
SyncDraftWorkflowResponse,
{
"result": "success",
"hash": workflow.unique_hash,
"updated_at": TimestampField().format(workflow.updated_at or workflow.created_at),
},
)
return {
"result": "success",
"hash": workflow.unique_hash,
"updated_at": TimestampField().format(workflow.updated_at or workflow.created_at),
}
@console_ns.route("/apps/<uuid:app_id>/advanced-chat/workflows/draft/run")
-14
View File
@@ -7,13 +7,10 @@ from configs import dify_config
from constants.languages import supported_language
from controllers.common.schema import query_params_from_model, register_schema_models
from controllers.console import console_ns
from controllers.console.auth.error import InvitationAccountMismatchError
from controllers.console.error import AccountInFreezeError, AlreadyActivateError
from extensions.ext_database import db
from libs.datetime_utils import naive_utc_now
from libs.helper import EmailStr, timezone
from libs.login import current_account_with_tenant
from libs.token import extract_access_token
from models import AccountStatus
from models.account import TenantAccountJoin, TenantAccountRole
from services.account_service import RegisterService, TenantService
@@ -139,12 +136,6 @@ class ActivateApi(Resource):
)
@console_ns.response(400, "Already activated or invalid token")
def post(self):
"""Accept an invitation without letting an existing session act for another account.
Token-only activation remains available for legacy clients. When the request already
carries a console session, that session must belong to the account encoded in the
invitation before the token is consumed or tenant membership is changed.
"""
args = ActivatePayload.model_validate(console_ns.payload)
normalized_request_email = args.email.lower() if args.email else None
@@ -155,11 +146,6 @@ class ActivateApi(Resource):
raise AlreadyActivateError()
account = invitation["account"]
if extract_access_token(request):
current_account, _ = current_account_with_tenant()
if current_account.id != account.id:
raise InvitationAccountMismatchError()
if dify_config.BILLING_ENABLED and BillingService.is_email_in_freeze(account.email):
raise AccountInFreezeError()
-6
View File
@@ -13,12 +13,6 @@ class InvalidEmailError(BaseHTTPException):
code = 400
class InvitationAccountMismatchError(BaseHTTPException):
error_code = "invitation_account_mismatch"
description = "This invitation was sent to another account. Please sign in with the invited account."
code = 403
class PasswordMismatchError(BaseHTTPException):
error_code = "password_mismatch"
description = "The passwords do not match."
+21 -35
View File
@@ -6,7 +6,6 @@ from flask import current_app, redirect, request
from flask_restx import Resource
from pydantic import BaseModel, Field
from werkzeug.exceptions import Unauthorized
from werkzeug.wrappers import Response
from configs import dify_config
from constants.languages import languages
@@ -128,20 +127,6 @@ def _preferred_interface_language(language: str | None = None) -> str:
return languages[0]
def _redirect_with_console_session(account: Account, target_url: str) -> Response:
"""Create a console session and attach its cookies to a redirect response."""
token_pair = AccountService.login(
account=account,
session=db.session(),
ip_address=extract_remote_ip(request),
)
response = redirect(target_url)
set_access_token_to_cookie(request, response, token_pair.access_token)
set_refresh_token_to_cookie(request, response, token_pair.refresh_token)
set_csrf_token_to_cookie(request, response, token_pair.csrf_token)
return response
@console_ns.route("/oauth/login/<provider>")
class OAuthLogin(Resource):
@console_ns.doc("oauth_login")
@@ -210,26 +195,16 @@ class OAuthCallback(Resource):
return redirect(f"{dify_config.CONSOLE_WEB_URL}/signin?message={urllib.parse.quote(str(e))}")
if invite_token and RegisterService.is_valid_invite_token(invite_token):
invitation = RegisterService.get_invitation_if_token_valid(
None,
None,
invite_token,
session=db.session(),
)
if not invitation:
return redirect(f"{dify_config.CONSOLE_WEB_URL}/signin?message=Invalid invitation token.")
if invitation["data"]["email"].lower() != user_info.email.lower():
message = "This invitation was sent to another account. Please sign in with the invited account."
query = urllib.parse.urlencode({"message": message, "invite_token": invite_token})
return redirect(f"{dify_config.CONSOLE_WEB_URL}/signin?{query}")
invitation = RegisterService.get_invitation_by_token(token=invite_token)
if invitation:
invitation_email = invitation.get("email", None)
invitation_email_normalized = (
invitation_email.lower() if isinstance(invitation_email, str) else invitation_email
)
if invitation_email_normalized != user_info.email.lower():
return redirect(f"{dify_config.CONSOLE_WEB_URL}/signin?message=Invalid invitation token.")
account = invitation["account"]
if account.status == AccountStatus.BANNED:
return redirect(f"{dify_config.CONSOLE_WEB_URL}/signin?message=Account is banned.")
AccountService.link_account_integrate(provider, user_info.id, account, session=db.session())
target_url = f"{dify_config.CONSOLE_WEB_URL}/signin/invite-settings?invite_token={invite_token}"
return _redirect_with_console_session(account, target_url)
return redirect(f"{dify_config.CONSOLE_WEB_URL}/signin/invite-settings?invite_token={invite_token}")
try:
account, oauth_new_user = _generate_account(provider, user_info, timezone=timezone, language=language)
@@ -264,10 +239,21 @@ class OAuthCallback(Resource):
"?message=Workspace not found, please contact system admin to invite you to join in a workspace."
)
token_pair = AccountService.login(
account=account,
session=db.session(),
ip_address=extract_remote_ip(request),
)
target_url = _get_redirect_target(redirect_url)
query_char = "&" if "?" in target_url else "?"
target_url = f"{target_url}{query_char}oauth_new_user={str(oauth_new_user).lower()}"
return _redirect_with_console_session(account, target_url)
response = redirect(target_url)
set_access_token_to_cookie(request, response, token_pair.access_token)
set_refresh_token_to_cookie(request, response, token_pair.refresh_token)
set_csrf_token_to_cookie(request, response, token_pair.csrf_token)
return response
def _get_account_by_openid_or_email(provider: str, user_info: OAuthUserInfo) -> Account | None:
+1 -8
View File
@@ -607,13 +607,6 @@ class DatasetListApi(Resource):
ReplaceMemberBindings(scope=RBACResourceWhitelistScope.ALL),
)
initialize_created_app_rbac_access_task.delay(current_tenant_id, current_user.id, dataset_id=dataset.id)
else:
enterprise_rbac_service.RBACService.DatasetAccess.replace_whitelist(
current_tenant_id,
current_user.id,
dataset.id,
ReplaceMemberBindings(scope=RBACResourceWhitelistScope.SPECIFIC),
)
permission_keys_map = enterprise_rbac_service.RBACService.DatasetPermissions.batch_get(
current_tenant_id,
@@ -882,7 +875,7 @@ class DatasetIndexingEstimateApi(Resource):
file_details = session.scalars(
select(UploadFile).where(UploadFile.tenant_id == current_tenant_id, UploadFile.id.in_(file_ids))
).all()
if not file_details:
if file_details is None:
raise NotFound("File not found.")
if file_details:
+1 -1
View File
@@ -645,7 +645,7 @@ class TrialChatAudioApi(TrialAppResource):
def post(self, current_user: Account, trial_app):
app_model = trial_app
file = request.files.get("file")
file = request.files["file"]
try:
# Get IDs before they might be detached from session
@@ -0,0 +1,5 @@
"""Typed Dify-owned KnowledgeFS Console product API."""
from . import resources
__all__ = ["resources"]
@@ -0,0 +1,56 @@
"""Stable, non-enumerating Console error contract for KnowledgeFS."""
from libs.exception import BaseHTTPException
class KnowledgeFSSpaceNotFoundHTTPError(BaseHTTPException):
error_code = "knowledge_fs_space_not_found"
description = "KnowledgeFS space was not found."
code = 404
class KnowledgeFSOperationUnavailableHTTPError(BaseHTTPException):
error_code = "knowledge_fs_operation_unavailable"
description = "KnowledgeFS operation is not available."
code = 503
class KnowledgeFSUpstreamUnavailableHTTPError(BaseHTTPException):
error_code = "knowledge_fs_upstream_unavailable"
description = "KnowledgeFS is unavailable."
code = 502
class KnowledgeFSInvalidRequestHTTPError(BaseHTTPException):
error_code = "knowledge_fs_invalid_request"
description = "KnowledgeFS request is invalid."
code = 400
class KnowledgeFSAccessDeniedHTTPError(BaseHTTPException):
error_code = "knowledge_fs_access_denied"
description = "KnowledgeFS operation is not allowed."
code = 403
class KnowledgeFSRateLimitHTTPError(BaseHTTPException):
error_code = "knowledge_fs_rate_limit_exceeded"
description = "KnowledgeFS operation rate limit exceeded."
code = 429
class KnowledgeFSQuotaExceededHTTPError(BaseHTTPException):
error_code = "knowledge_fs_quota_exceeded"
description = "KnowledgeFS operation quota exceeded."
code = 403
__all__ = [
"KnowledgeFSAccessDeniedHTTPError",
"KnowledgeFSInvalidRequestHTTPError",
"KnowledgeFSOperationUnavailableHTTPError",
"KnowledgeFSQuotaExceededHTTPError",
"KnowledgeFSRateLimitHTTPError",
"KnowledgeFSSpaceNotFoundHTTPError",
"KnowledgeFSUpstreamUnavailableHTTPError",
]
File diff suppressed because it is too large Load Diff
@@ -1,419 +0,0 @@
"""Authenticated transport adapter for the Console-to-KnowledgeFS proxy.
These raw Blueprint routes deliberately stay outside Dify's OpenAPI surface:
KnowledgeFS owns the wire contract consumed by the frontend. The catch-all path
avoids resource-specific Dify controllers, while the forwarding module consumes
only the operations explicitly enabled by Dify's product registry. The registry
can be validated explicitly against the pinned KnowledgeFS contract during development.
Console auth and contract-specific dataset RBAC run before forwarding. Request
bodies are capped at 64 MiB, JSON and binary responses have separate bounds,
SSE responses remain streaming with a bounded idle read timeout, and only safe
response headers are exposed. Operation-specific upstream error mappings are
applied before Console JSON error handling; the default maps 401 to 502 so it
cannot trigger browser-session recovery and preserves resource-level 403.
"""
from __future__ import annotations
import logging
from collections.abc import Callable, Iterator
from functools import wraps
from http import HTTPStatus
from typing import NoReturn, cast
import httpx
from flask import Response, request, stream_with_context
from flask.typing import ResponseReturnValue
from werkzeug.exceptions import (
BadGateway,
Forbidden,
GatewayTimeout,
HTTPException,
NotFound,
RequestEntityTooLarge,
ServiceUnavailable,
default_exceptions,
)
from configs import dify_config
from controllers.console import api, bp
from controllers.console.wraps import (
account_initialization_required,
cloud_edition_billing_rate_limit_check,
setup_required,
)
from core.helper import ssrf_proxy
from libs.login import current_account_with_tenant, login_required
from services.knowledge_fs_operations import KnowledgeFSMethod
from services.knowledge_fs_proxy import (
KnowledgeFSAccessDeniedError,
KnowledgeFSAuthorization,
KnowledgeFSConfigurationError,
KnowledgeFSRouteNotAllowedError,
KnowledgeFSTimeoutError,
KnowledgeFSTransportError,
KnowledgeFSUpstreamResponse,
authorize_knowledge_fs_request,
get_knowledge_fs_operation,
proxy_authorized_knowledge_fs_request,
proxy_knowledge_fs_request,
)
logger = logging.getLogger(__name__)
type _KnowledgeFSRequestForwarder = Callable[
[str | None, str | None, bytes | None, bytes | None],
KnowledgeFSUpstreamResponse,
]
_MAX_PROXY_BODY_BYTES = 64 * 1024 * 1024
_RESPONSE_HEADER_ALLOWLIST = (
"Cache-Control",
"Content-Disposition",
"Content-Type",
"Retry-After",
"X-Trace-Id",
)
_RESPONSE_HEADER_DENYLIST = frozenset(
{
"authorization",
"connection",
"cookie",
"keep-alive",
"proxy-authenticate",
"proxy-authorization",
"set-cookie",
"te",
"trailer",
"transfer-encoding",
"upgrade",
}
)
def _console_api_errors[**P](
view: Callable[P, ResponseReturnValue],
) -> Callable[P, ResponseReturnValue]:
"""Route raw Blueprint exceptions through the Console API JSON handlers."""
@wraps(view)
def decorated(*args: P.args, **kwargs: P.kwargs) -> ResponseReturnValue:
try:
return view(*args, **kwargs)
except Exception as exc:
return api.handle_error(exc)
return decorated
def _knowledge_fs_enabled[**P](
view: Callable[P, ResponseReturnValue],
) -> Callable[P, ResponseReturnValue]:
"""Hide the complete KnowledgeFS route surface while the bridge is disabled."""
@wraps(view)
def decorated(*args: P.args, **kwargs: P.kwargs) -> ResponseReturnValue:
if not dify_config.KNOWLEDGE_FS_ENABLED:
raise NotFound()
return view(*args, **kwargs)
return decorated
def _translate_proxy_error(exc: Exception, *, tenant_id: str) -> NoReturn:
"""Map forwarding failures to the stable Console HTTP error surface."""
if isinstance(exc, KnowledgeFSRouteNotAllowedError):
raise NotFound() from exc
if isinstance(exc, KnowledgeFSAccessDeniedError):
raise Forbidden() from exc
if isinstance(exc, KnowledgeFSConfigurationError):
logger.error("KnowledgeFS request was blocked by invalid configuration for tenant_id=%s", tenant_id)
raise ServiceUnavailable("KnowledgeFS integration is misconfigured") from exc
if isinstance(exc, KnowledgeFSTimeoutError):
raise GatewayTimeout("KnowledgeFS request timed out") from exc
if isinstance(exc, KnowledgeFSTransportError):
logger.warning("KnowledgeFS transport request failed for tenant_id=%s", tenant_id)
raise BadGateway("KnowledgeFS is unavailable") from exc
raise exc
def _knowledge_fs_operation_access_required(
view: Callable[[KnowledgeFSAuthorization], ResponseReturnValue],
) -> Callable[[KnowledgeFSMethod, str], ResponseReturnValue]:
"""Authorize one declared operation before billing and request-body work."""
@wraps(view)
def decorated(method: KnowledgeFSMethod, upstream_path: str) -> ResponseReturnValue:
current_user, tenant_id = current_account_with_tenant()
try:
authorization = authorize_knowledge_fs_request(
account=current_user,
tenant_id=tenant_id,
method=method,
path=upstream_path,
)
except KnowledgeFSRouteNotAllowedError as exc:
raise NotFound() from exc
except KnowledgeFSAccessDeniedError as exc:
_translate_proxy_error(exc, tenant_id=tenant_id)
return view(authorization)
return decorated
def _request_body() -> bytes:
"""Read the raw body up to the proxy limit or raise RequestEntityTooLarge."""
body = request.stream.read(_MAX_PROXY_BODY_BYTES + 1)
if len(body) > _MAX_PROXY_BODY_BYTES:
raise RequestEntityTooLarge("KnowledgeFS proxy request body is too large")
return body
def _stream_response_body(
upstream: httpx.Response,
*,
tenant_id: str,
max_response_bytes: int,
) -> Iterator[bytes]:
"""Yield one bounded SSE response and always release its pooled connection."""
total_bytes = 0
try:
for chunk in upstream.iter_bytes():
total_bytes += len(chunk)
if total_bytes > max_response_bytes:
logger.warning("KnowledgeFS stream exceeded the proxy limit for tenant_id=%s", tenant_id)
raise ssrf_proxy.ResponseTooLargeError(f"response exceeded {max_response_bytes} bytes")
yield chunk
finally:
upstream.close()
def _proxy_response(
upstream_result: KnowledgeFSUpstreamResponse,
*,
tenant_id: str,
contract_response_headers: tuple[str, ...],
max_response_bytes: int,
) -> Response:
"""Expose raw content, status, and allowlisted headers from KnowledgeFS.
Raises:
HTTPException: KnowledgeFS returns a status normalized by the operation contract.
"""
upstream = upstream_result.response
mapped_status = dict(upstream_result.operation.error_status_map).get(upstream.status_code)
if mapped_status is not None:
upstream.close()
description = "KnowledgeFS upstream request failed"
if upstream.status_code == HTTPStatus.UNAUTHORIZED:
description = "KnowledgeFS authentication failed"
logger.error(
"KnowledgeFS rejected the Dify server credential with HTTP %s for tenant_id=%s",
upstream.status_code,
tenant_id,
)
exception_type = default_exceptions.get(mapped_status)
if exception_type is None:
exception = HTTPException(description)
exception.code = mapped_status
raise exception
raise exception_type(description)
allowed_header_names = dict.fromkeys(
name.lower() for name in (*_RESPONSE_HEADER_ALLOWLIST, *contract_response_headers)
)
headers = {
name: value
for name in allowed_header_names
if name not in _RESPONSE_HEADER_DENYLIST
if (value := upstream.headers.get(name)) is not None
}
if upstream_result.response_kind == "stream":
response = Response(
stream_with_context( # pyrefly: ignore[no-matching-overload]
_stream_response_body(
upstream,
tenant_id=tenant_id,
max_response_bytes=max_response_bytes,
)
),
status=upstream.status_code,
headers=headers,
)
response.call_on_close(upstream.close)
return response
try:
content = upstream.content
finally:
upstream.close()
return Response(content, status=upstream.status_code, headers=headers)
def _proxy_current_request(
*,
method: KnowledgeFSMethod,
tenant_id: str,
forward: _KnowledgeFSRequestForwarder,
) -> Response:
"""Forward the current raw request through one preconfigured service entry."""
if not dify_config.KNOWLEDGE_FS_ENABLED:
raise NotFound()
try:
proxy_result = forward(
request.headers.get("Accept"),
request.content_type,
request.query_string or None,
_request_body() if method != "GET" else None,
)
except (
KnowledgeFSConfigurationError,
KnowledgeFSAccessDeniedError,
KnowledgeFSRouteNotAllowedError,
KnowledgeFSTimeoutError,
KnowledgeFSTransportError,
) as exc:
_translate_proxy_error(exc, tenant_id=tenant_id)
return _proxy_response(
proxy_result,
tenant_id=tenant_id,
contract_response_headers=proxy_result.operation.response_headers,
max_response_bytes=proxy_result.operation.max_response_bytes,
)
def _proxy_request(
method: KnowledgeFSMethod,
upstream_path: str,
) -> Response:
"""Authorize and forward the current request through the combined service use case."""
if not dify_config.KNOWLEDGE_FS_ENABLED:
raise NotFound()
current_user, tenant_id = current_account_with_tenant()
def forward(
accept: str | None,
content_type: str | None,
query: bytes | None,
body: bytes | None,
) -> KnowledgeFSUpstreamResponse:
return proxy_knowledge_fs_request(
account=current_user,
method=method,
path=upstream_path,
tenant_id=tenant_id,
accept=accept,
content_type=content_type,
query=query,
body=body,
request_headers=request.headers,
)
return _proxy_current_request(method=method, tenant_id=tenant_id, forward=forward)
def _proxy_authorized_request(authorization: KnowledgeFSAuthorization) -> Response:
"""Forward the current request using one previously authorized operation capability.
Args:
authorization: Request-scoped capability produced before billing and body parsing.
Returns:
The filtered response returned by KnowledgeFS.
Raises:
HTTPException: The integration is disabled or forwarding fails.
"""
operation = authorization.operation
tenant_id = authorization.tenant_id
def forward(
accept: str | None,
content_type: str | None,
query: bytes | None,
body: bytes | None,
) -> KnowledgeFSUpstreamResponse:
return proxy_authorized_knowledge_fs_request(
authorization=authorization,
accept=accept,
content_type=content_type,
query=query,
body=body,
request_headers=request.headers,
)
return _proxy_current_request(method=operation.method, tenant_id=tenant_id, forward=forward)
@_knowledge_fs_enabled
@_knowledge_fs_operation_access_required
@cloud_edition_billing_rate_limit_check("knowledge")
def _proxy_knowledge_fs_non_get(
authorization: KnowledgeFSAuthorization,
) -> ResponseReturnValue:
"""Apply knowledge billing checks to one allowlisted non-GET operation."""
return _proxy_authorized_request(authorization)
@bp.route(
"/knowledge-fs/<path:upstream_path>",
methods=["OPTIONS"],
provide_automatic_options=False,
)
@_console_api_errors
@_knowledge_fs_enabled
def proxy_knowledge_fs_options(upstream_path: str) -> ResponseReturnValue:
"""Complete a CORS preflight only for an enabled Console operation."""
requested_method = cast(KnowledgeFSMethod, request.headers.get("Access-Control-Request-Method", "").upper())
try:
get_knowledge_fs_operation(requested_method, upstream_path)
except KnowledgeFSRouteNotAllowedError as exc:
raise NotFound() from exc
return Response(status=HTTPStatus.NO_CONTENT)
@bp.route(
"/knowledge-fs/<path:upstream_path>",
methods=["GET"],
provide_automatic_options=False,
)
@_console_api_errors
@_knowledge_fs_enabled
@setup_required
@login_required
@account_initialization_required
def proxy_knowledge_fs_get(upstream_path: str) -> ResponseReturnValue:
"""Forward one authenticated, dataset-readable GET request.
Args:
upstream_path: Relative KFS path captured after the Console proxy prefix.
Returns:
The filtered raw KnowledgeFS response or a Console JSON error response.
"""
if request.method != "GET":
raise NotFound()
return _proxy_request("GET", upstream_path)
@bp.route(
"/knowledge-fs/<path:upstream_path>",
methods=["DELETE", "PATCH", "POST", "PUT"],
provide_automatic_options=False,
)
@_console_api_errors
@_knowledge_fs_enabled
@setup_required
@login_required
@account_initialization_required
def proxy_knowledge_fs_write(upstream_path: str) -> ResponseReturnValue:
"""Forward one authenticated non-GET request under its contract access policy.
Args:
upstream_path: Relative KFS path captured after the Console proxy prefix.
Returns:
The filtered raw KnowledgeFS response or a Console JSON error response.
"""
method = cast(KnowledgeFSMethod, request.method)
return _proxy_knowledge_fs_non_get(method, upstream_path)
-106
View File
@@ -1,106 +0,0 @@
"""Console onboarding APIs.
This module keeps Step-by-step Tour persistence account-scoped. Workspace IDs
are accepted only as presentation overrides; UI-only state such as minimized
panels or the currently active task stays on the frontend. PATCH requests are
action-based so callers do not replace server-side arrays with stale snapshots.
"""
from datetime import datetime
from typing import Literal, cast
from flask_restx import Resource
from pydantic import BaseModel, ConfigDict, Field, model_validator
from controllers.common.schema import register_response_schema_models, register_schema_models
from extensions.ext_database import db
from fields.base import ResponseModel
from libs.helper import dump_response
from libs.login import login_required
from models import Account
from services.step_by_step_tour_service import StepByStepTourPatch, StepByStepTourService
from . import console_ns
from .wraps import account_initialization_required, setup_required, with_current_tenant_id, with_current_user
StepByStepTourAction = Literal[
"skip",
"complete_task",
"uncomplete_task",
"enable_current_workspace",
"disable_current_workspace",
]
StepByStepTourTaskId = Literal["home", "studio", "knowledge", "integration"]
class StepByStepTourStatePatchPayload(BaseModel):
action: StepByStepTourAction = Field(description="State update action")
task_id: StepByStepTourTaskId | None = Field(default=None, description="Task ID for task actions")
model_config = ConfigDict(extra="forbid")
@model_validator(mode="after")
def validate_patch_shape(self) -> "StepByStepTourStatePatchPayload":
task_actions = {"complete_task", "uncomplete_task"}
if self.action in task_actions and self.task_id is None:
raise ValueError("task_id is required for task actions")
if self.action not in task_actions and self.task_id is not None:
raise ValueError("task_id is only supported for task actions")
return self
class StepByStepTourStateResponse(ResponseModel):
first_workspace_id: str | None = None
skipped: bool = False
completed_task_ids: list[StepByStepTourTaskId] = Field(default_factory=list)
manually_enabled_workspace_ids: list[str] = Field(default_factory=list)
manually_disabled_workspace_ids: list[str] = Field(default_factory=list)
updated_at: datetime | None = None
register_schema_models(console_ns, StepByStepTourStatePatchPayload)
register_response_schema_models(console_ns, StepByStepTourStateResponse)
@console_ns.route("/onboarding/step-by-step-tour/state")
class StepByStepTourStateApi(Resource):
@console_ns.doc("get_step_by_step_tour_state")
@console_ns.doc(description="Get account-level Step-by-step Tour state")
@console_ns.response(200, "Success", console_ns.models[StepByStepTourStateResponse.__name__])
@setup_required
@login_required
@account_initialization_required
@with_current_user
@with_current_tenant_id
def get(self, current_tenant_id: str, current_user: Account):
return dump_response(
StepByStepTourStateResponse,
StepByStepTourService.get_state(
account=current_user,
current_tenant_id=current_tenant_id,
session=db.session,
),
)
@console_ns.doc("patch_step_by_step_tour_state")
@console_ns.doc(description="Update account-level Step-by-step Tour state")
@console_ns.expect(console_ns.models[StepByStepTourStatePatchPayload.__name__])
@console_ns.response(200, "Success", console_ns.models[StepByStepTourStateResponse.__name__])
@setup_required
@login_required
@account_initialization_required
@with_current_user
@with_current_tenant_id
def patch(self, current_tenant_id: str, current_user: Account):
payload = StepByStepTourStatePatchPayload.model_validate(console_ns.payload or {})
patch = cast(StepByStepTourPatch, payload.model_dump(exclude_unset=True, exclude_none=True))
return dump_response(
StepByStepTourStateResponse,
StepByStepTourService.patch_state(
account=current_user,
current_tenant_id=current_tenant_id,
patch=patch,
session=db.session,
),
)
+27 -10
View File
@@ -4,16 +4,20 @@ from http import HTTPStatus
from flask import redirect
from flask_restx import Resource
from pydantic import BaseModel, Field
from werkzeug.exceptions import Conflict, Forbidden, NotFound
from werkzeug.exceptions import Conflict, NotFound
from controllers.common.fields import RedirectResponse
from controllers.common.schema import register_response_schema_models, register_schema_models
from controllers.console import console_ns
from controllers.console.wraps import (
RBACPermission,
RBACResourceScope,
account_initialization_required,
cloud_edition_billing_enabled,
cloud_edition_billing_paid_plan_required,
is_admin_or_owner_required,
only_edition_cloud,
rbac_permission_required,
setup_required,
)
from extensions.ext_database import db
@@ -21,7 +25,6 @@ from fields.base import ResponseModel
from libs.archive_storage import get_export_storage
from libs.helper import dump_response
from libs.login import current_account_with_tenant, login_required
from models import TenantAccountRole
from services.retention.workflow_run.archive_download_preparation import ARCHIVE_DOWNLOAD_MIME_TYPE
from services.retention.workflow_run.archive_download_task_cache import (
WorkflowRunArchiveDownloadStatus,
@@ -95,13 +98,11 @@ register_response_schema_models(
)
def _current_owner_or_admin_ids() -> tuple[str, str]:
"""Return current Cloud workspace IDs for an owner or admin, independently of enterprise RBAC."""
def _current_ids() -> tuple[str, str]:
"""Return current `(tenant_id, account_id)` or raise when no workspace is selected."""
current_user, current_tenant_id = current_account_with_tenant()
if not current_tenant_id:
raise NotFound("Current workspace not found")
if not TenantAccountRole.is_privileged_role(current_user.current_role):
raise Forbidden()
return current_tenant_id, current_user.id
@@ -123,8 +124,12 @@ class WorkflowRunArchivesApi(Resource):
@only_edition_cloud
@cloud_edition_billing_enabled
@cloud_edition_billing_paid_plan_required
@is_admin_or_owner_required
@rbac_permission_required(
RBACResourceScope.WORKSPACE, RBACPermission.WORKSPACE_ROLE_MANAGE, resource_required=False
)
def get(self):
tenant_id, _ = _current_owner_or_admin_ids()
tenant_id, _ = _current_ids()
return dump_response(WorkflowRunArchiveListResponse, list_workflow_run_archives(db.session(), tenant_id))
@@ -144,8 +149,12 @@ class WorkflowRunArchiveDownloadsApi(Resource):
@only_edition_cloud
@cloud_edition_billing_enabled
@cloud_edition_billing_paid_plan_required
@is_admin_or_owner_required
@rbac_permission_required(
RBACResourceScope.WORKSPACE, RBACPermission.WORKSPACE_ROLE_MANAGE, resource_required=False
)
def post(self):
tenant_id, account_id = _current_owner_or_admin_ids()
tenant_id, account_id = _current_ids()
payload = WorkflowRunArchiveDownloadPayload.model_validate(console_ns.payload or {})
try:
task = create_workflow_run_archive_download_task(
@@ -171,8 +180,12 @@ class WorkflowRunArchiveDownloadApi(Resource):
@only_edition_cloud
@cloud_edition_billing_enabled
@cloud_edition_billing_paid_plan_required
@is_admin_or_owner_required
@rbac_permission_required(
RBACResourceScope.WORKSPACE, RBACPermission.WORKSPACE_ROLE_MANAGE, resource_required=False
)
def get(self, download_id: str):
tenant_id, _ = _current_owner_or_admin_ids()
tenant_id, _ = _current_ids()
try:
task = get_workflow_run_archive_download_task(tenant_id=tenant_id, download_id=download_id)
except WorkflowRunArchiveDownloadTaskNotFoundError as exc:
@@ -196,8 +209,12 @@ class WorkflowRunArchiveDownloadFileApi(Resource):
@only_edition_cloud
@cloud_edition_billing_enabled
@cloud_edition_billing_paid_plan_required
@is_admin_or_owner_required
@rbac_permission_required(
RBACResourceScope.WORKSPACE, RBACPermission.WORKSPACE_ROLE_MANAGE, resource_required=False
)
def get(self, download_id: str):
tenant_id, _ = _current_owner_or_admin_ids()
tenant_id, _ = _current_ids()
try:
task = get_ready_workflow_run_archive_download_task(tenant_id=tenant_id, download_id=download_id)
except WorkflowRunArchiveDownloadTaskNotFoundError as exc:
@@ -67,15 +67,6 @@ from services.plugin.plugin_parameter_service import PluginParameterService
from services.plugin.plugin_permission_service import PluginPermissionService
from services.tools.tools_transform_service import ToolTransformService
_PLUGIN_PACKAGE_UPLOAD_PARAMS = {
"pkg": {
"description": "Plugin package to upload",
"in": "formData",
"type": "file",
"required": True,
}
}
class AutoUpgradeSettingsResponse(TypedDict):
strategy_setting: TenantPluginAutoUpgradeStrategySetting
@@ -654,7 +645,6 @@ class PluginAssetApi(Resource):
@console_ns.route("/workspaces/current/plugin/upload/pkg")
class PluginUploadFromPkgApi(Resource):
@console_ns.doc(consumes=["multipart/form-data"], params=_PLUGIN_PACKAGE_UPLOAD_PARAMS)
@console_ns.response(200, "Success", console_ns.models[PluginDecodeResponse.__name__])
@setup_required
@login_required
+14 -3
View File
@@ -346,7 +346,13 @@ class RBACRoleItemApi(Resource):
def put(self, role_id):
tenant_id, account_id = _current_ids()
request = _payload(_RoleUpsertRequest)
role = svc.RBACService.Roles.update(tenant_id, account_id, str(role_id), request.to_mutation())
role = svc.RBACService.KnowledgeFSRoleMutations.update_role(
tenant_id,
account_id,
str(role_id),
request.to_mutation(),
session=db.session(),
)
return _dump(role)
@login_required
@@ -356,7 +362,12 @@ class RBACRoleItemApi(Resource):
@console_ns.response(200, "Success", console_ns.models[svc.RBACRole.__name__])
def delete(self, role_id):
tenant_id, account_id = _current_ids()
svc.RBACService.Roles.delete(tenant_id, account_id, str(role_id))
svc.RBACService.KnowledgeFSRoleMutations.delete_role(
tenant_id,
account_id,
str(role_id),
session=db.session(),
)
return {"result": "success"}
@@ -915,7 +926,7 @@ class RBACMemberRolesApi(Resource):
tenant_id, account_id = _current_ids()
request = _payload(_ReplaceMemberRolesRequest)
return _dump(
svc.RBACService.MemberRoles.replace(
svc.RBACService.KnowledgeFSRoleMutations.replace_member_roles(
tenant_id,
account_id,
str(member_id),
@@ -8,6 +8,7 @@ from controllers.inner_api.plugin.wraps import get_user_tenant, plugin_data
from controllers.inner_api.wraps import plugin_inner_api_only
from core.plugin.backwards_invocation.app import PluginAppBackwardsInvocation
from core.plugin.backwards_invocation.base import BaseBackwardsInvocationResponse
from core.plugin.backwards_invocation.datasource import PluginDatasourceBackwardsInvocation
from core.plugin.backwards_invocation.encrypt import PluginEncrypter
from core.plugin.backwards_invocation.model import PluginModelBackwardsInvocation
from core.plugin.backwards_invocation.node import PluginNodeBackwardsInvocation
@@ -15,10 +16,12 @@ from core.plugin.backwards_invocation.tool import PluginToolBackwardsInvocation
from core.plugin.entities.request import (
RequestFetchAppInfo,
RequestInvokeApp,
RequestInvokeDatasource,
RequestInvokeEncrypt,
RequestInvokeLLM,
RequestInvokeLLMWithStructuredOutput,
RequestInvokeModeration,
RequestInvokeMultimodalEmbedding,
RequestInvokeParameterExtractorNode,
RequestInvokeQuestionClassifierNode,
RequestInvokeRerank,
@@ -27,6 +30,7 @@ from core.plugin.entities.request import (
RequestInvokeTextEmbedding,
RequestInvokeTool,
RequestInvokeTTS,
RequestListModels,
RequestRequestDownloadFile,
RequestRequestUploadFile,
)
@@ -118,6 +122,36 @@ class PluginInvokeTextEmbeddingApi(Resource):
return jsonable_encoder(BaseBackwardsInvocationResponse(error=str(e)))
@inner_api_ns.route("/invoke/multimodal-embedding")
class PluginInvokeMultimodalEmbeddingApi(Resource):
@get_user_tenant
@setup_required
@plugin_inner_api_only
@plugin_data(payload_type=RequestInvokeMultimodalEmbedding)
@inner_api_ns.doc("plugin_invoke_multimodal_embedding")
@inner_api_ns.doc(description="Invoke multimodal embedding models through Dify model management")
@inner_api_ns.doc(
responses={
200: "Multimodal embedding successful",
401: "Unauthorized - invalid API key",
404: "Service not available",
}
)
def post(self, user_model: Account | EndUser, tenant_model: Tenant, payload: RequestInvokeMultimodalEmbedding):
try:
return jsonable_encoder(
BaseBackwardsInvocationResponse(
data=PluginModelBackwardsInvocation.invoke_multimodal_embedding(
user_id=user_model.id,
tenant=tenant_model,
payload=payload,
)
)
)
except Exception as e:
return jsonable_encoder(BaseBackwardsInvocationResponse(error=str(e)))
@inner_api_ns.route("/invoke/rerank")
class PluginInvokeRerankApi(Resource):
@get_user_tenant
@@ -144,6 +178,70 @@ class PluginInvokeRerankApi(Resource):
return jsonable_encoder(BaseBackwardsInvocationResponse(error=str(e)))
@inner_api_ns.route("/invoke/model-catalog")
class PluginModelCatalogApi(Resource):
@get_user_tenant
@setup_required
@plugin_inner_api_only
@plugin_data(payload_type=RequestListModels)
@inner_api_ns.doc("plugin_model_catalog")
@inner_api_ns.doc(description="List tenant-active models managed by Dify")
@inner_api_ns.doc(
responses={
200: "Model catalog lookup successful",
401: "Unauthorized - invalid API key",
404: "Service not available",
}
)
def post(self, user_model: Account | EndUser, tenant_model: Tenant, payload: RequestListModels):
try:
return jsonable_encoder(
BaseBackwardsInvocationResponse(
data=PluginModelBackwardsInvocation.list_models(
tenant_id=tenant_model.id,
user_id=user_model.id,
payload=payload,
)
)
)
except Exception as e:
return jsonable_encoder(BaseBackwardsInvocationResponse(error=str(e)))
@inner_api_ns.route("/invoke/datasource")
class PluginInvokeDatasourceApi(Resource):
"""Invoke an installed datasource with credentials resolved inside Dify."""
@get_user_tenant
@setup_required
@plugin_inner_api_only
@plugin_data(payload_type=RequestInvokeDatasource)
@inner_api_ns.doc("plugin_invoke_datasource")
@inner_api_ns.doc(description="Invoke datasource plugins through Dify credential management")
@inner_api_ns.doc(
responses={
200: "Datasource invocation successful (streaming response)",
401: "Unauthorized - invalid API key",
404: "Datasource provider, datasource, or credential not found",
}
)
def post(
self,
user_model: Account | EndUser,
tenant_model: Tenant,
payload: RequestInvokeDatasource,
):
response = PluginDatasourceBackwardsInvocation.invoke(
user_id=user_model.id,
tenant=tenant_model,
payload=payload,
)
return length_prefixed_response(
0xF,
PluginDatasourceBackwardsInvocation.convert_to_event_stream(response),
)
@inner_api_ns.route("/invoke/tts")
class PluginInvokeTTSApi(Resource):
@get_user_tenant
+2
View File
@@ -38,6 +38,7 @@ from .dataset import (
)
from .dataset.rag_pipeline import rag_pipeline_workflow
from .end_user import end_user
from .knowledge_fs import resources as knowledge_fs_resources
from .workspace import models
__all__ = [
@@ -54,6 +55,7 @@ __all__ = [
"hit_testing",
"human_input_form",
"index",
"knowledge_fs_resources",
"message",
"metadata",
"models",
+1 -1
View File
@@ -101,7 +101,7 @@ class AudioApi(Resource):
Accepts an audio file upload and returns the transcribed text.
"""
file = request.files.get("file")
file = request.files["file"]
try:
response = AudioService.transcript_asr(
+8 -16
View File
@@ -531,22 +531,14 @@ class DatasetListApi(DatasetApiResource):
except services.errors.dataset.DatasetNameDuplicateError:
raise DatasetNameDuplicateError()
if dify_config.RBAC_ENABLED:
if payload.permission == DatasetPermissionEnum.ALL_TEAM:
RBACService.DatasetAccess.replace_whitelist(
tenant_id,
current_user.id,
dataset.id,
ReplaceMemberBindings(scope=RBACResourceWhitelistScope.ALL),
)
initialize_created_app_rbac_access_task.delay(tenant_id, current_user.id, dataset_id=dataset.id)
else:
RBACService.DatasetAccess.replace_whitelist(
tenant_id,
current_user.id,
dataset.id,
ReplaceMemberBindings(scope=RBACResourceWhitelistScope.SPECIFIC),
)
if payload.permission == DatasetPermissionEnum.ALL_TEAM and dify_config.RBAC_ENABLED:
RBACService.DatasetAccess.replace_whitelist(
tenant_id,
current_user.id,
dataset.id,
ReplaceMemberBindings(scope=RBACResourceWhitelistScope.ALL),
)
initialize_created_app_rbac_access_task.delay(tenant_id, current_user.id, dataset_id=dataset.id)
return _dump_service_dataset_detail(dataset, session=session), 200
@@ -0,0 +1,5 @@
"""KnowledgeFS-specific Service API authenticated by resource credentials."""
from . import resources
__all__ = ["resources"]
@@ -0,0 +1,49 @@
"""Stable KnowledgeFS Service API error contract."""
from libs.exception import BaseHTTPException
class KnowledgeFSInvalidCredentialHTTPError(BaseHTTPException):
error_code = "knowledge_fs_invalid_credential"
description = "Invalid KnowledgeFS service credential."
code = 401
class KnowledgeFSServiceOperationUnavailableHTTPError(BaseHTTPException):
error_code = "knowledge_fs_operation_unavailable"
description = "KnowledgeFS operation is not available."
code = 503
class KnowledgeFSServiceUpstreamUnavailableHTTPError(BaseHTTPException):
error_code = "knowledge_fs_upstream_unavailable"
description = "KnowledgeFS is unavailable."
code = 502
class KnowledgeFSServiceInvalidRequestHTTPError(BaseHTTPException):
error_code = "knowledge_fs_invalid_request"
description = "KnowledgeFS request is invalid."
code = 400
class KnowledgeFSServiceRateLimitHTTPError(BaseHTTPException):
error_code = "knowledge_fs_rate_limit_exceeded"
description = "KnowledgeFS operation rate limit exceeded."
code = 429
class KnowledgeFSServiceQuotaExceededHTTPError(BaseHTTPException):
error_code = "knowledge_fs_quota_exceeded"
description = "KnowledgeFS operation quota exceeded."
code = 403
__all__ = [
"KnowledgeFSInvalidCredentialHTTPError",
"KnowledgeFSServiceInvalidRequestHTTPError",
"KnowledgeFSServiceOperationUnavailableHTTPError",
"KnowledgeFSServiceQuotaExceededHTTPError",
"KnowledgeFSServiceRateLimitHTTPError",
"KnowledgeFSServiceUpstreamUnavailableHTTPError",
]
File diff suppressed because it is too large Load Diff
+1 -1
View File
@@ -76,7 +76,7 @@ class AudioApi(WebApiResource):
@web_ns.response(200, "Success", web_ns.models[AudioToTextResponse.__name__])
def post(self, app_model: App, end_user: EndUser):
"""Convert audio to text"""
file = request.files.get("file")
file = request.files["file"]
try:
response = AudioService.transcript_asr(
+7 -17
View File
@@ -1,6 +1,6 @@
from typing import Any, Self
from pydantic import AliasChoices, Field
from pydantic import AliasChoices, Field, computed_field
from sqlalchemy import select
from werkzeug.exceptions import Forbidden
@@ -9,13 +9,11 @@ from controllers.common.schema import register_response_schema_models
from controllers.web import web_ns
from controllers.web.wraps import WebApiResource
from extensions.ext_database import db
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, EndUser, IconType, Site
from models.model import App, EndUser, Site
from services.feature_service import FeatureModel, FeatureService
from services.file_service import FileService
class WebSiteResponse(ResponseModel):
@@ -34,7 +32,11 @@ class WebSiteResponse(ResponseModel):
prompt_public: bool | None = None
show_workflow_steps: bool | None = None
use_icon_as_answer_icon: bool | None = None
icon_url: 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)
class WebModelConfigResponse(ResponseModel):
@@ -86,7 +88,6 @@ class WebAppSiteResponse(ResponseModel):
end_user_id: str | None,
features: FeatureModel,
can_replace_logo: bool,
icon_url: str | None = None,
) -> Self:
custom_config = None
if can_replace_logo:
@@ -101,7 +102,6 @@ class WebAppSiteResponse(ResponseModel):
)
site_response = WebSiteResponse.model_validate(site, from_attributes=True)
site_response.icon_url = icon_url if icon_url is not None else build_icon_url(site.icon_type, site.icon)
if features.billing.enabled and not features.webapp_copyright_enabled:
site_response.copyright = None
site_response.input_placeholder = None
@@ -123,15 +123,6 @@ register_response_schema_models(
)
def _build_site_icon_url(*, site: Site, tenant_id: str) -> str | None:
"""Use direct S3 URLs only in Cloud Mode and preserve preview URLs elsewhere."""
if site.icon_type != IconType.IMAGE or not site.icon:
return None
if dify_config.EDITION == "CLOUD" and StorageType(dify_config.STORAGE_TYPE) == StorageType.S3:
return FileService(db.engine).get_file_presigned_url(file_id=site.icon, tenant_id=tenant_id)
return build_icon_url(site.icon_type, site.icon)
@web_ns.route("/site")
class AppSiteApi(WebApiResource):
@web_ns.doc("Get App Site Info")
@@ -168,5 +159,4 @@ class AppSiteApi(WebApiResource):
end_user_id=end_user.id,
features=features,
can_replace_logo=features.can_replace_logo,
icon_url=_build_site_icon_url(site=site, tenant_id=tenant.id),
).model_dump(mode="json")
+12
View File
@@ -14,7 +14,10 @@ from core.app.apps.base_app_queue_manager import AppQueueManager
from core.app.apps.base_app_runner import AppRunner
from core.app.entities.app_invoke_entities import (
AgentChatAppGenerateEntity,
DifyRunContext,
InvokeFrom,
ModelConfigWithCredentialsEntity,
UserFrom,
)
from core.app.file_access import DatabaseFileAccessController
from core.callback_handler.agent_tool_callback_handler import DifyAgentCallbackHandler
@@ -148,6 +151,15 @@ class BaseAgentRunner(AppRunner):
user_id=self.user_id,
invoke_from=self.application_generate_entity.invoke_from,
)
invoke_from = self.application_generate_entity.invoke_from
if isinstance(invoke_from, InvokeFrom):
tool_entity.runtime.dify_run_context = DifyRunContext(
tenant_id=self.tenant_id,
app_id=self.app_config.app_id,
user_id=self.user_id,
user_from=(UserFrom.ACCOUNT if invoke_from.runs_as_account() else UserFrom.END_USER),
invoke_from=invoke_from,
)
assert tool_entity.entity.description
message_tool = PromptMessageTool(
name=tool.tool_name,
@@ -625,11 +625,7 @@ class AdvancedChatAppGenerator(MessageBasedAppGenerator):
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,
),
draft_var_saver_factory=self._get_draft_var_saver_factory(invoke_from, account=user),
)
return AdvancedChatAppGenerateResponseConverter.convert(response=response, invoke_from=invoke_from)
+29 -14
View File
@@ -682,27 +682,42 @@ class AgentAppGenerator(MessageBasedAppGenerator):
if draft_type == AgentConfigDraftType.DEBUG_BUILD.value
else AgentConfigDraftType.DRAFT
)
if effective_draft_type == AgentConfigDraftType.DRAFT:
from services.agent.composer_service import AgentComposerService
return AgentComposerService.get_or_create_normal_agent_draft(
session=session,
tenant_id=tenant_id,
agent=agent,
created_by=agent.updated_by or agent.created_by,
)
if not account_id:
raise AgentAppGeneratorError("Build draft requires an account user")
stmt = select(AgentConfigDraft).where(
AgentConfigDraft.tenant_id == tenant_id,
AgentConfigDraft.agent_id == agent.id,
AgentConfigDraft.draft_type == AgentConfigDraftType.DEBUG_BUILD,
AgentConfigDraft.account_id == account_id,
AgentConfigDraft.draft_type == effective_draft_type,
)
if effective_draft_type == AgentConfigDraftType.DEBUG_BUILD:
if not account_id:
raise AgentAppGeneratorError("Build draft requires an account user")
stmt = stmt.where(AgentConfigDraft.account_id == account_id)
else:
stmt = stmt.where(AgentConfigDraft.account_id.is_(None))
draft = session.scalar(stmt.order_by(AgentConfigDraft.updated_at.desc()).limit(1))
if draft is not None:
return draft
raise AgentAppGeneratorError("Agent build draft not found")
if effective_draft_type == AgentConfigDraftType.DEBUG_BUILD:
raise AgentAppGeneratorError("Agent build draft not found")
_, snapshot, agent_soul = AgentAppGenerator._resolve_agent_by_id(
tenant_id=tenant_id,
agent_id=agent.id,
snapshot_id=agent.active_config_snapshot_id,
session=session,
)
draft = AgentConfigDraft(
tenant_id=tenant_id,
agent_id=agent.id,
draft_type=AgentConfigDraftType.DRAFT,
account_id=None,
draft_owner_key="",
base_snapshot_id=snapshot.id,
config_snapshot=agent_soul,
created_by=agent.created_by,
updated_by=agent.updated_by,
)
session.add(draft)
session.flush()
return draft
@staticmethod
def _resolve_agent_by_id(
@@ -159,6 +159,7 @@ class AgentAppRuntimeRequestBuilder:
user_from=cast(DifyExecutionContextUserFrom, context.dify_context.user_from.value),
invoke_from=cast(DifyExecutionContextInvokeFrom, context.dify_context.invoke_from.value),
agent_mode="agent_app",
trace_id=context.dify_context.trace_session_id,
),
# ENG-616: expand slash-menu mention tokens to canonical names so
# no frontend-internal {{#…#}} marker ever reaches the model.
+1 -10
View File
@@ -32,7 +32,6 @@ class _DebuggerDraftVariableSaver:
self,
*,
account: Account,
tenant_id: str,
app_id: str,
node_id: str,
node_type: NodeType,
@@ -40,7 +39,6 @@ class _DebuggerDraftVariableSaver:
enclosing_node_id: str | None = None,
) -> None:
self._account = account
self._tenant_id = tenant_id
self._app_id = app_id
self._node_id = node_id
self._node_type = node_type
@@ -51,7 +49,6 @@ class _DebuggerDraftVariableSaver:
with Session(db.engine) as session, session.begin():
DraftVariableSaverImpl(
session=session,
tenant_id=self._tenant_id,
app_id=self._app_id,
node_id=self._node_id,
node_type=self._node_type,
@@ -290,12 +287,7 @@ class BaseAppGenerator:
@final
@staticmethod
def _get_draft_var_saver_factory(
invoke_from: InvokeFrom,
account: Account | EndUser,
*,
tenant_id: str,
) -> DraftVariableSaverFactory:
def _get_draft_var_saver_factory(invoke_from: InvokeFrom, account: Account | EndUser) -> DraftVariableSaverFactory:
if invoke_from == InvokeFrom.DEBUGGER:
assert isinstance(account, Account)
@@ -308,7 +300,6 @@ class BaseAppGenerator:
) -> DraftVariableSaver:
return _DebuggerDraftVariableSaver(
account=account,
tenant_id=tenant_id,
app_id=app_id,
node_id=node_id,
node_type=node_type,
@@ -349,7 +349,6 @@ class PipelineGenerator(BaseAppGenerator):
draft_var_saver_factory = self._get_draft_var_saver_factory(
invoke_from,
user,
tenant_id=pipeline.tenant_id,
)
# return response or stream generator
response = self._handle_response(
+1 -5
View File
@@ -399,11 +399,7 @@ class WorkflowAppGenerator(BaseAppGenerator):
worker_thread.start()
draft_var_saver_factory = self._get_draft_var_saver_factory(
invoke_from,
user,
tenant_id=app_model.tenant_id,
)
draft_var_saver_factory = self._get_draft_var_saver_factory(invoke_from, user)
# return response or stream generator
response = self._handle_response(
@@ -38,6 +38,11 @@ class DatasourcePluginProviderController(ABC):
):
raise ToolProviderCredentialValidationError("Invalid credentials")
def validate_credentials(self, user_id: str, credentials: dict[str, Any]) -> None:
"""Validate credential shape and value against this installed provider declaration."""
self.validate_credentials_format(credentials)
self._validate_credentials(user_id, credentials)
@property
def provider_type(self) -> DatasourceProviderType:
"""
@@ -0,0 +1,110 @@
"""Datasource backward invocation through Dify's tenant-bound runtime.
This module is the internal service boundary used by trusted callers such as
KnowledgeFS. It resolves installed provider declarations and Dify-owned
credential references before reaching ``PluginDatasourceManager``; callers
must never provide raw datasource credentials or a plugin-daemon API key.
"""
from collections.abc import Generator
from typing import Any, cast
from pydantic import BaseModel
from core.datasource.datasource_manager import DatasourceManager
from core.datasource.entities.datasource_entities import (
OnlineDriveBrowseFilesRequest,
OnlineDriveDownloadFileRequest,
)
from core.datasource.online_document.online_document_plugin import OnlineDocumentDatasourcePlugin
from core.datasource.online_drive.online_drive_plugin import OnlineDriveDatasourcePlugin
from core.datasource.website_crawl.website_crawl_plugin import WebsiteCrawlDatasourcePlugin
from core.plugin.backwards_invocation.base import BaseBackwardsInvocation
from core.plugin.entities.request import RequestInvokeDatasource
from models.account import Tenant
from models.provider_ids import DatasourceProviderID
from services.datasource_provider_service import DatasourceProviderService
class PluginDatasourceBackwardsInvocation(BaseBackwardsInvocation):
"""Resolve and invoke a datasource without exposing credential material to the caller."""
@classmethod
def invoke(
cls,
*,
user_id: str,
tenant: Tenant,
payload: RequestInvokeDatasource,
) -> Generator[BaseModel | dict[str, Any], None, None]:
"""Yield datasource messages for one validated inner-runtime request."""
provider_id = DatasourceProviderID(payload.provider)
canonical_provider_id = str(provider_id)
controller = DatasourceManager.get_datasource_plugin_provider(
provider_id=canonical_provider_id,
tenant_id=tenant.id,
datasource_type=payload.datasource_type,
)
if controller.entity.provider_type != payload.datasource_type:
raise ValueError("Datasource provider type mismatch")
# Resolving the datasource from the installed declaration prevents a caller
# from dispatching an arbitrary datasource name under a valid plugin ID.
runtime = controller.get_datasource(payload.datasource)
credentials = DatasourceProviderService().get_datasource_credentials(
tenant_id=tenant.id,
provider=provider_id.provider_name,
plugin_id=provider_id.plugin_id,
credential_id=payload.credential_id,
)
if controller.need_credentials and not credentials:
raise ValueError("Datasource credential not found")
if payload.operation == "validate_credentials":
controller.validate_credentials(user_id=user_id, credentials=credentials)
yield {"result": True}
return
runtime.runtime.credentials = credentials
provider_type = runtime.datasource_provider_type()
match payload.operation:
case "get_website_crawl":
website = cast(WebsiteCrawlDatasourcePlugin, runtime)
yield from website.get_website_crawl(
user_id=user_id,
datasource_parameters=payload.datasource_parameters,
provider_type=provider_type,
)
case "get_online_document_pages":
document = cast(OnlineDocumentDatasourcePlugin, runtime)
yield from document.get_online_document_pages(
user_id=user_id,
datasource_parameters=payload.datasource_parameters,
provider_type=provider_type,
)
case "get_online_document_page_content":
if payload.page is None:
raise ValueError("Online-document page input is required")
document = cast(OnlineDocumentDatasourcePlugin, runtime)
yield from document.get_online_document_page_content(
user_id=user_id,
datasource_parameters=payload.page,
provider_type=provider_type,
)
case "online_drive_browse_files":
drive = cast(OnlineDriveDatasourcePlugin, runtime)
yield from drive.online_drive_browse_files(
user_id=user_id,
request=OnlineDriveBrowseFilesRequest.model_validate(payload.request),
provider_type=provider_type,
)
case "online_drive_download_file":
drive = cast(OnlineDriveDatasourcePlugin, runtime)
yield from drive.online_drive_download_file(
user_id=user_id,
request=OnlineDriveDownloadFileRequest.model_validate(payload.request),
provider_type=provider_type,
)
case _:
raise ValueError(f"Unsupported datasource operation: {payload.operation}")
+109 -2
View File
@@ -1,22 +1,31 @@
import tempfile
from binascii import hexlify, unhexlify
from collections.abc import Generator
from collections.abc import Generator, Mapping
from enum import Enum
from typing import Any
from pydantic import BaseModel
from core.app.llm import deduct_llm_quota
from core.llm_generator.output_parser.structured_output import invoke_llm_with_structured_output
from core.model_manager import ModelManager
from core.plugin.backwards_invocation.base import BaseBackwardsInvocation
from core.plugin.entities.request import (
InvokableModelCatalogItem,
InvokableModelCatalogPage,
RequestInvokeLLM,
RequestInvokeLLMWithStructuredOutput,
RequestInvokeModeration,
RequestInvokeMultimodalEmbedding,
RequestInvokeRerank,
RequestInvokeSpeech2Text,
RequestInvokeSummary,
RequestInvokeTextEmbedding,
RequestInvokeTTS,
RequestListModels,
)
from core.plugin.impl.model_runtime_factory import create_plugin_provider_manager
from core.plugin.plugin_service import PluginService
from core.tools.entities.tool_entities import ToolProviderType
from core.tools.utils.model_invocation_utils import ModelInvocationUtils
from graphon.model_runtime.entities.llm_entities import (
@@ -33,6 +42,20 @@ from graphon.model_runtime.entities.message_entities import (
)
from graphon.model_runtime.entities.model_entities import ModelType
from models.account import Tenant
from models.provider_ids import ModelProviderID
def _json_compatible(value: Any) -> Any:
"""Convert model-runtime metadata into stable JSON-compatible values."""
if isinstance(value, BaseModel):
return value.model_dump(mode="json")
if isinstance(value, Enum):
return value.value
if isinstance(value, Mapping):
return {str(_json_compatible(key)): _json_compatible(child) for key, child in value.items()}
if isinstance(value, list | tuple | set):
return [_json_compatible(child) for child in value]
return value
class PluginModelBackwardsInvocation(BaseBackwardsInvocation):
@@ -183,7 +206,30 @@ class PluginModelBackwardsInvocation(BaseBackwardsInvocation):
)
# invoke model
response = model_instance.invoke_text_embedding(texts=payload.texts)
response = model_instance.invoke_text_embedding(texts=payload.texts, input_type=payload.input_type)
return response
@classmethod
def invoke_multimodal_embedding(
cls,
user_id: str,
tenant: Tenant,
payload: RequestInvokeMultimodalEmbedding,
):
"""Invoke multimodal embedding through the tenant-bound model instance."""
model_instance = cls._get_bound_model_instance(
tenant_id=tenant.id,
user_id=user_id,
provider=payload.provider,
model_type=payload.model_type,
model=payload.model,
)
response = model_instance.invoke_multimodal_embedding(
multimodel_documents=[document.model_dump(exclude_none=True) for document in payload.documents],
input_type=payload.input_type,
)
return response
@@ -210,6 +256,67 @@ class PluginModelBackwardsInvocation(BaseBackwardsInvocation):
return response
@classmethod
def list_models(
cls,
tenant_id: str,
user_id: str,
payload: RequestListModels,
) -> InvokableModelCatalogPage:
"""List only models that are active for the tenant's Dify configuration."""
provider_manager = create_plugin_provider_manager(tenant_id=tenant_id, user_id=user_id)
active_models = provider_manager.get_configurations(tenant_id).get_models(
model_type=payload.model_type,
only_active=True,
)
installed_identities: dict[str, str] = {}
for plugin in PluginService.list(tenant_id):
existing = installed_identities.get(plugin.plugin_id)
if existing is not None and existing != plugin.plugin_unique_identifier:
raise ValueError(f"Ambiguous installed identity for model plugin {plugin.plugin_id}")
installed_identities[plugin.plugin_id] = plugin.plugin_unique_identifier
requested_provider = str(ModelProviderID(payload.provider)) if payload.provider else None
matched_models = [
model
for model in active_models
if (requested_provider is None or model.provider.provider == requested_provider)
and (payload.model is None or model.model == payload.model)
]
matched_models.sort(key=lambda model: (model.provider.provider, model.model))
page_models = matched_models[payload.offset : payload.offset + payload.limit]
items: list[InvokableModelCatalogItem] = []
for model in page_models:
provider_id = ModelProviderID(model.provider.provider)
unique_identifier = installed_identities.get(provider_id.plugin_id)
if unique_identifier is None:
raise ValueError(f"Installed identity not found for active model plugin {provider_id.plugin_id}")
items.append(
InvokableModelCatalogItem(
plugin_id=provider_id.plugin_id,
plugin_unique_identifier=unique_identifier,
provider=provider_id.provider_name,
model=model.model,
model_type=model.model_type,
capabilities={
"deprecated": model.deprecated,
"features": _json_compatible(model.features or []),
"fetchFrom": _json_compatible(model.fetch_from),
"modelProperties": _json_compatible(model.model_properties),
"modelType": model.model_type.value,
"status": _json_compatible(model.status),
},
)
)
next_offset = payload.offset + len(page_models)
return InvokableModelCatalogPage(
items=items,
next_offset=next_offset if next_offset < len(matched_models) else None,
)
@classmethod
def invoke_tts(cls, user_id: str, tenant: Tenant, payload: RequestInvokeTTS):
"""
+113 -2
View File
@@ -6,6 +6,11 @@ from typing import Any, Literal
from flask import Response
from pydantic import BaseModel, ConfigDict, Field, field_validator, model_validator
from core.datasource.entities.datasource_entities import (
DatasourceProviderType,
GetOnlineDocumentPageContentRequest,
)
from core.entities.embedding_type import EmbeddingInputType
from core.entities.provider_entities import BasicProviderConfig
from core.plugin.utils.http_parser import deserialize_response
from core.workflow.file_reference import is_canonical_file_reference
@@ -54,6 +59,61 @@ class RequestInvokeTool(BaseModel):
credential_id: str | None = None
DatasourceInvocationOperation = Literal[
"get_online_document_page_content",
"get_online_document_pages",
"get_website_crawl",
"online_drive_browse_files",
"online_drive_download_file",
"validate_credentials",
]
class RequestInvokeDatasource(BaseModel):
"""Invoke one installed datasource using a Dify-owned credential reference.
Raw credentials are intentionally not part of this contract. ``tenant_id`` and
``user_id`` are consumed by the inner-API request context, while the remaining
fields select an installed provider declaration and an operation-specific input.
"""
tenant_id: str = Field(min_length=1, max_length=512)
user_id: str = Field(min_length=1, max_length=512)
provider: str = Field(min_length=1, max_length=768)
datasource: str = Field(min_length=1, max_length=256)
datasource_type: DatasourceProviderType
credential_id: str = Field(min_length=1, max_length=512)
operation: DatasourceInvocationOperation
datasource_parameters: dict[str, Any] = Field(default_factory=dict)
page: GetOnlineDocumentPageContentRequest | None = None
request: dict[str, Any] | None = None
model_config = ConfigDict(extra="forbid")
@model_validator(mode="after")
def validate_operation_payload(self) -> "RequestInvokeDatasource":
expected_type = {
"get_online_document_page_content": DatasourceProviderType.ONLINE_DOCUMENT,
"get_online_document_pages": DatasourceProviderType.ONLINE_DOCUMENT,
"get_website_crawl": DatasourceProviderType.WEBSITE_CRAWL,
"online_drive_browse_files": DatasourceProviderType.ONLINE_DRIVE,
"online_drive_download_file": DatasourceProviderType.ONLINE_DRIVE,
"validate_credentials": self.datasource_type,
}[self.operation]
if self.datasource_type != expected_type:
raise ValueError(f"{self.operation} requires datasource_type {expected_type.value}")
page_required = self.operation == "get_online_document_page_content"
if page_required != (self.page is not None):
raise ValueError("page is required only for get_online_document_page_content")
request_required = self.operation in {"online_drive_browse_files", "online_drive_download_file"}
if request_required != (self.request is not None):
raise ValueError("request is required only for online-drive operations")
return self
class BaseRequestInvokeModel(BaseModel):
provider: str
model: str
@@ -115,6 +175,25 @@ class RequestInvokeTextEmbedding(BaseRequestInvokeModel):
model_type: ModelType = ModelType.TEXT_EMBEDDING
texts: list[str]
input_type: EmbeddingInputType = EmbeddingInputType.DOCUMENT
class MultimodalEmbeddingDocument(BaseModel):
"""A document accepted by a multimodal text-embedding model."""
content: str
content_type: str
file_id: str | None = None
model_config = ConfigDict(extra="forbid")
class RequestInvokeMultimodalEmbedding(BaseRequestInvokeModel):
"""Request to invoke a multimodal text-embedding model."""
model_type: ModelType = ModelType.TEXT_EMBEDDING
documents: list[MultimodalEmbeddingDocument] = Field(min_length=1)
input_type: EmbeddingInputType = EmbeddingInputType.DOCUMENT
class RequestInvokeRerank(BaseRequestInvokeModel):
@@ -125,8 +204,40 @@ class RequestInvokeRerank(BaseRequestInvokeModel):
model_type: ModelType = ModelType.RERANK
query: str
docs: list[str]
score_threshold: float
top_n: int
score_threshold: float | None = None
top_n: int | None = None
class RequestListModels(BaseModel):
"""Tenant-scoped query for models that Dify can invoke."""
model_type: Literal[ModelType.LLM, ModelType.TEXT_EMBEDDING, ModelType.RERANK]
provider: str | None = None
model: str | None = None
offset: int = Field(default=0, ge=0)
limit: int = Field(default=50, ge=1, le=100)
model_config = ConfigDict(protected_namespaces=())
class InvokableModelCatalogItem(BaseModel):
"""Installed identity and active Dify capability metadata for one model."""
plugin_id: str
plugin_unique_identifier: str
provider: str
model: str
model_type: ModelType
capabilities: dict[str, Any] = Field(default_factory=dict)
model_config = ConfigDict(protected_namespaces=())
class InvokableModelCatalogPage(BaseModel):
"""Offset page returned by the internal model catalog endpoint."""
items: list[InvokableModelCatalogItem] = Field(default_factory=list)
next_offset: int | None = None
class RequestInvokeTTS(BaseRequestInvokeModel):
+1 -14
View File
@@ -1,7 +1,7 @@
import inspect
import json
import logging
from collections.abc import Callable, Generator, Mapping
from collections.abc import Callable, Generator
from typing import Any, cast
from urllib.parse import unquote
@@ -23,7 +23,6 @@ from core.plugin.impl.exc import (
PluginLLMPollingUnsupportedError,
PluginNotFoundError,
PluginPermissionDeniedError,
PluginRuntimeError,
PluginUniqueIdentifierError,
)
from core.trigger.errors import (
@@ -376,18 +375,6 @@ class BasePluginClient:
# type `PluginLLMPollingUnsupportedError`.
case PluginLLMPollingUnsupportedError.__name__:
raise PluginLLMPollingUnsupportedError(description=error_object.get("message"))
case PluginRuntimeError.__name__:
args = error_object.get("args")
lambda_request_id = args.get("request_id") if isinstance(args, Mapping) else None
if not isinstance(lambda_request_id, str):
lambda_request_id = None
runtime_message = error_object.get("message")
if not isinstance(runtime_message, str):
runtime_message = "Plugin runtime request failed"
raise PluginRuntimeError(
description=runtime_message,
lambda_request_id=lambda_request_id,
)
case _:
raise PluginInvokeError(description=message)
case PluginDaemonInternalServerError.__name__:
-12
View File
@@ -49,18 +49,6 @@ class PluginDaemonBadRequestError(PluginDaemonClientSideError):
description: str = "Bad Request"
class PluginRuntimeError(PluginDaemonInternalError):
"""A plugin runtime failed before it could return a valid plugin response."""
lambda_request_id: str | None
def __init__(self, description: str, lambda_request_id: str | None = None) -> None:
self.lambda_request_id = lambda_request_id
if lambda_request_id:
description = description.replace(f"RequestId: {lambda_request_id} Error: ", "", 1)
super().__init__(description)
class PluginInvokeError(PluginDaemonClientSideError, ValueError):
description: str = "Invoke Error"
+1 -4
View File
@@ -131,10 +131,7 @@ class WaterCrawlAPIClient(BaseAPIClient):
content_type = response.headers.get("Content-Type", "")
media_type = content_type.split(";", 1)[0].strip().lower()
if media_type == "application/json":
try:
return response.json() or {}
except ValueError as exc:
raise ValueError("Invalid JSON response from WaterCrawl") from exc
return response.json() or {}
if media_type == "application/octet-stream":
return response.content
@@ -1,15 +1,5 @@
"""WaterCrawl domain exceptions.
These exceptions are constructed from upstream HTTP responses, which may be
JSON API errors or plain text/HTML proxy errors. Keep the exception type stable
even when the body is not JSON so callers can handle WaterCrawl failures by
domain type instead of low-level parser errors.
"""
import json
from typing import Any, override
from httpx import Response
from typing import override
class WaterCrawlError(Exception):
@@ -17,16 +7,11 @@ class WaterCrawlError(Exception):
class WaterCrawlBadRequestError(WaterCrawlError):
def __init__(self, response: Response):
def __init__(self, response):
self.status_code = response.status_code
self.response = response
try:
data: Any = response.json()
except ValueError:
data = {}
if not isinstance(data, dict):
data = {}
self.message = data.get("message") or response.text or "Unknown error occurred"
data = response.json()
self.message = data.get("message", "Unknown error occurred")
self.errors = data.get("errors", {})
super().__init__(self.message)
+10
View File
@@ -10,6 +10,7 @@ class RBACResourceScope(StrEnum):
APP = "app"
DATASET = "dataset"
KNOWLEDGE_FS = "knowledge_space"
WORKSPACE = "workspace"
@@ -57,6 +58,15 @@ class RBACPermission(StrEnum):
DATASET_EXTERNAL_CONNECT = "dataset_external_connect"
DATASET_IMPORT_EXPORT_DSL = "dataset_import_export_dsl"
KNOWLEDGE_FS_READ = "knowledge_space_read"
KNOWLEDGE_FS_CREATE = "knowledge_space_create"
KNOWLEDGE_FS_EDIT = "knowledge_space_edit"
KNOWLEDGE_FS_DELETE = "knowledge_space_delete"
KNOWLEDGE_FS_ACCESS_CONFIG = "knowledge_space_access_config"
KNOWLEDGE_FS_API_KEY_MANAGE = "knowledge_space_api_key_manage"
KNOWLEDGE_FS_DOCUMENT_WRITE = "knowledge_space_document_write"
KNOWLEDGE_FS_QUERY = "knowledge_space_query"
WORKSPACE_MEMBER_MANAGE = "workspace_member_manage"
WORKSPACE_ROLE_MANAGE = "workspace_role_manage"
API_EXTENSION_MANAGE = "api_extension_manage"
+15 -28
View File
@@ -58,6 +58,7 @@ class Tool(ABC):
if self.runtime and self.runtime.runtime_parameters:
tool_parameters.update(self.runtime.runtime_parameters)
# try parse tool parameters into the correct type
tool_parameters = self._transform_tool_parameters_type(tool_parameters)
result = self._invoke(
@@ -86,14 +87,14 @@ class Tool(ABC):
return result
def _transform_tool_parameters_type(self, tool_parameters: dict[str, Any]) -> dict[str, Any]:
"""Transform declared tool parameter values without resolving runtime schemas."""
"""
Transform tool parameters type
"""
# Temp fix for the issue that the tool parameters will be converted to empty while validating the credentials
result = deepcopy(tool_parameters)
for parameter in self.entity.parameters or []:
if parameter.name in tool_parameters:
if parameter.multiple:
result[parameter.name] = parameter.init_frontend_parameter(result.get(parameter.name))
else:
result[parameter.name] = parameter.type.cast_value(tool_parameters[parameter.name])
result[parameter.name] = parameter.type.cast_value(tool_parameters[parameter.name])
return result
@@ -195,31 +196,17 @@ class Tool(ABC):
}:
continue
is_multiple_select = parameter.multiple and parameter.type in {
ToolParameter.ToolParameterType.SELECT,
ToolParameter.ToolParameterType.DYNAMIC_SELECT,
}
if is_multiple_select:
item_schema: dict[str, Any] = {"type": "string"}
if parameter.type == ToolParameter.ToolParameterType.SELECT and parameter.options:
item_schema["enum"] = [option.value for option in parameter.options]
parameter_schema: dict[str, Any] = {"type": "array", "items": item_schema}
else:
parameter_schema = (
{
"type": parameter.type.as_normal_type(),
"description": parameter.llm_description or "",
}
if parameter.input_schema is None
else deepcopy(parameter.input_schema)
)
parameter_schema: dict[str, Any] = (
{
"type": parameter.type.as_normal_type(),
"description": parameter.llm_description or "",
}
if parameter.input_schema is None
else deepcopy(parameter.input_schema)
)
parameter_schema.setdefault("description", parameter.llm_description or "")
if (
not is_multiple_select
and parameter.type == ToolParameter.ToolParameterType.SELECT
and parameter.options
):
if parameter.type == ToolParameter.ToolParameterType.SELECT and parameter.options:
parameter_schema["enum"] = [option.value for option in parameter.options]
schema["properties"][parameter.name] = parameter_schema
+2 -1
View File
@@ -2,7 +2,7 @@ from typing import Any
from pydantic import BaseModel, Field
from core.app.entities.app_invoke_entities import InvokeFrom
from core.app.entities.app_invoke_entities import DifyRunContext, InvokeFrom
from core.plugin.entities.plugin_daemon import CredentialType
from core.tools.entities.tool_entities import ToolInvokeFrom
@@ -20,6 +20,7 @@ class ToolRuntime(BaseModel):
tool_id: str | None = None
invoke_from: InvokeFrom | None = None
tool_invoke_from: ToolInvokeFrom | None = None
dify_run_context: DifyRunContext | None = Field(default=None, exclude=True, repr=False)
credentials: dict[str, Any] = Field(default_factory=dict)
credential_type: CredentialType = Field(default=CredentialType.API_KEY)
runtime_parameters: dict[str, Any] = Field(default_factory=dict)
@@ -1,4 +1,5 @@
- audio
- code
- knowledge_fs
- time
- webscraper
@@ -0,0 +1,6 @@
<svg xmlns="http://www.w3.org/2000/svg" viewBox="0 0 64 64" fill="none">
<rect width="64" height="64" rx="14" fill="#155EEF"/>
<path d="M17 17h20c5.5 0 10 4.5 10 10v20H27c-5.5 0-10-4.5-10-10V17Z" fill="white" fill-opacity=".96"/>
<path d="M26 27h12M26 34h12M26 41h7" stroke="#155EEF" stroke-width="4" stroke-linecap="round"/>
</svg>

After

Width:  |  Height:  |  Size: 340 B

@@ -0,0 +1,9 @@
from typing import Any, override
from core.tools.builtin_tool.provider import BuiltinToolProviderController
class KnowledgeFSProvider(BuiltinToolProviderController):
@override
def _validate_credentials(self, user_id: str, credentials: dict[str, Any]) -> None:
_ = (user_id, credentials)
@@ -0,0 +1,13 @@
identity:
author: Dify
name: knowledge_fs
label:
en_US: KnowledgeFS
zh_Hans: KnowledgeFS
description:
en_US: Run explicitly bound KnowledgeFS operations from an Agent or Workflow.
zh_Hans: 从 Agent 或 Workflow 执行显式绑定的 KnowledgeFS 操作。
icon: icon.svg
tags:
- rag
@@ -0,0 +1,60 @@
from collections.abc import Generator
from typing import Any, override
from pydantic import ValidationError
from sqlalchemy.orm import Session, sessionmaker
from core.tools.builtin_tool.tool import BuiltinTool
from core.tools.entities.tool_entities import ToolInvokeFrom, ToolInvokeMessage
from core.tools.errors import ToolInvokeError
from models.knowledge_fs import KnowledgeFSAppSpaceJoinType
from services.knowledge_fs.app_execution_capability import KnowledgeResourceRef
from services.knowledge_fs.product_dto import KnowledgeFSResearchTaskCreatePayload
from services.knowledge_fs.runtime import create_knowledge_fs_runtime
class KnowledgeFSCreateResearchTaskTool(BuiltinTool):
@override
def _invoke(
self,
session: Session,
user_id: str,
tool_parameters: dict[str, Any],
conversation_id: str | None = None,
app_id: str | None = None,
message_id: str | None = None,
) -> Generator[ToolInvokeMessage, None, None]:
_ = (user_id, conversation_id, app_id, message_id)
run_context = self.runtime.dify_run_context
if run_context is None or self.runtime.tenant_id != run_context.tenant_id:
raise ToolInvokeError("KnowledgeFS requires a trusted Dify run context")
match self.runtime.tool_invoke_from:
case ToolInvokeFrom.AGENT:
caller_kind = KnowledgeFSAppSpaceJoinType.AGENT
case ToolInvokeFrom.WORKFLOW:
caller_kind = KnowledgeFSAppSpaceJoinType.WORKFLOW
case _:
raise ToolInvokeError("KnowledgeFS is only available to Agent and Workflow callers")
try:
resource = KnowledgeResourceRef.model_validate(tool_parameters.get("resource"))
payload_data: dict[str, object] = {
"query": tool_parameters.get("query"),
}
mode = tool_parameters.get("mode")
if mode:
payload_data["mode"] = mode
payload = KnowledgeFSResearchTaskCreatePayload.model_validate(payload_data)
runtime = create_knowledge_fs_runtime(sessionmaker(bind=session.get_bind(), expire_on_commit=False))
response = runtime.app_capabilities.create_research_task(
run_context=run_context,
caller_kind=caller_kind,
resource=resource,
payload=payload,
)
except ToolInvokeError:
raise
except (ValidationError, RuntimeError, ValueError) as exc:
raise ToolInvokeError(str(exc)) from exc
yield self.create_json_message(response.model_dump(mode="json", by_alias=True))
@@ -0,0 +1,75 @@
identity:
name: create_research_task
author: Dify
label:
en_US: Create KnowledgeFS Research Task
zh_Hans: 创建 KnowledgeFS Research 任务
description:
human:
en_US: Create a Research task in an explicitly bound KnowledgeFS space.
zh_Hans: 在显式绑定的 KnowledgeFS 空间中创建 Research 任务。
llm: Create a durable research task using an explicitly configured KnowledgeFS resource.
parameters:
- name: resource
type: object
required: true
label:
en_US: KnowledgeFS resource
zh_Hans: KnowledgeFS 资源
human_description:
en_US: A typed KnowledgeFS control-space reference configured by the app author.
zh_Hans: 由应用作者配置的类型化 KnowledgeFS control-space 引用。
form: form
input_schema:
type: object
additionalProperties: false
properties:
kind:
type: string
const: knowledge_fs
control_space_id:
type: string
minLength: 1
required:
- kind
- control_space_id
- name: query
type: string
required: true
label:
en_US: Research query
zh_Hans: Research 查询
human_description:
en_US: The question the Research task should investigate.
zh_Hans: Research 任务需要调查的问题。
llm_description: The question to investigate with KnowledgeFS.
form: llm
- name: mode
type: select
required: false
label:
en_US: Mode
zh_Hans: 模式
human_description:
en_US: Optional KnowledgeFS retrieval mode.
zh_Hans: 可选的 KnowledgeFS 检索模式。
llm_description: Optional retrieval mode. Use auto unless the task needs a specific mode.
form: llm
options:
- value: auto
label:
en_US: Auto
zh_Hans: 自动
- value: fast
label:
en_US: Fast
zh_Hans: 快速
- value: deep
label:
en_US: Deep
zh_Hans: 深度
- value: research
label:
en_US: Research
zh_Hans: Research
+5 -36
View File
@@ -292,7 +292,9 @@ class ToolInvokeMessageBinary(BaseModel):
class ToolParameter(PluginParameter):
"""Tool-specific parameter declaration and invocation-value normalization."""
"""
Overrides type
"""
class ToolParameterType(StrEnum):
"""
@@ -331,28 +333,12 @@ class ToolParameter(PluginParameter):
LLM = auto() # will be set by LLM
type: ToolParameterType = Field(..., description="The type of the parameter")
multiple: bool = Field(
default=False,
description="Whether the parameter is multiple select, only valid for select or dynamic-select type",
)
human_description: I18nObject | None = Field(default=None, description="The description presented to the user")
form: ToolParameterForm = Field(..., description="The form of the parameter, schema/form/llm")
llm_description: str | None = None
# MCP object and array type parameters use this field to store the schema
input_schema: dict[str, Any] | None = None
@model_validator(mode="after")
def validate_multiple(self) -> ToolParameter:
supports_multiple = self.type in {
self.ToolParameterType.SELECT,
self.ToolParameterType.DYNAMIC_SELECT,
}
if self.multiple and not supports_multiple:
raise ValueError("multiple is only valid for select and dynamic-select parameters")
if supports_multiple and self.default is not None and (isinstance(self.default, list) != self.multiple):
raise ValueError("default must be a list exactly when multiple is true")
return self
@classmethod
def get_simple_instance(
cls,
@@ -392,25 +378,8 @@ class ToolParameter(PluginParameter):
options=option_objs,
)
def init_frontend_parameter(self, value: Any) -> Any:
"""Normalize a value against this tool parameter's full declaration."""
if not self.multiple:
return init_frontend_parameter(self, self.type, value)
parameter_value = self.default if value is None else value
if parameter_value is None:
parameter_value = []
if not isinstance(parameter_value, list):
raise ValueError(f"tool parameter {self.name} must be a list when multiple is true")
if not all(isinstance(item, str) for item in parameter_value):
raise ValueError(f"tool parameter {self.name} must contain only strings")
if self.required and not parameter_value:
raise ValueError(f"tool parameter {self.name} not found in tool config")
if self.type == self.ToolParameterType.SELECT:
options = [option.value for option in self.options]
if any(item not in options for item in parameter_value):
raise ValueError(f"tool parameter {self.name} value {parameter_value} not in options {options}")
return parameter_value
def init_frontend_parameter(self, value: Any):
return init_frontend_parameter(self, self.type, value)
class ToolProviderIdentity(BaseModel):
+1
View File
@@ -493,6 +493,7 @@ class DifyToolNodeRuntime(ToolNodeRuntimeProtocol):
self._run_context.invoke_from,
variable_pool,
)
tool_runtime.runtime.dify_run_context = self._run_context
except ToolNodeError:
raise
except Exception as exc:
@@ -2,7 +2,7 @@ from __future__ import annotations
from collections.abc import Mapping
from dataclasses import dataclass
from typing import Any, Final, Literal, Protocol
from typing import Any, Literal, Protocol, cast
from dify_agent.layers.dify_core_tools import DifyCoreToolConfig, DifyCoreToolProviderType, DifyCoreToolsLayerConfig
from dify_agent.layers.dify_plugin import (
@@ -29,13 +29,6 @@ from models.provider_ids import ToolProviderID
from models.tools import WorkflowToolProvider
from services.tools.mcp_tools_manage_service import MCPToolManageService
_CORE_TOOL_PROVIDER_TYPES: Final[dict[ToolProviderType, DifyCoreToolProviderType]] = {
ToolProviderType.BUILT_IN: "builtin",
ToolProviderType.API: "api",
ToolProviderType.WORKFLOW: "workflow",
ToolProviderType.MCP: "mcp",
}
class WorkflowAgentDifyToolsBuildError(ValueError):
"""Raised when Agent Soul tools cannot be prepared for Agent backend."""
@@ -242,7 +235,7 @@ class WorkflowAgentDifyToolsBuilder:
if tool_config.tool_name is not None:
expanded.append(tool_config)
continue
provider_type = tool_config.provider_type
provider_type = ToolProviderType.value_of(tool_config.provider_type)
provider_id = self._provider_id(tool_config)
try:
tool_names = self._provider_declared_tool_names(
@@ -286,7 +279,7 @@ class WorkflowAgentDifyToolsBuilder:
tenant_id: str,
tool_config: AgentSoulDifyToolConfig,
) -> AgentSoulDifyToolConfig:
if tool_config.provider_type is not ToolProviderType.MCP:
if tool_config.provider_type != ToolProviderType.MCP.value:
return tool_config
provider_id = self._mcp_provider_id_resolver(tenant_id=tenant_id, provider_id=self._provider_id(tool_config))
return tool_config.model_copy(update={"provider_id": provider_id, "plugin_id": None, "provider": None})
@@ -333,7 +326,7 @@ class WorkflowAgentDifyToolsBuilder:
def _to_agent_tool_entity(tool_config: AgentSoulDifyToolConfig) -> AgentToolEntity:
assert tool_config.tool_name is not None
return AgentToolEntity(
provider_type=tool_config.provider_type,
provider_type=ToolProviderType.value_of(tool_config.provider_type),
provider_id=WorkflowAgentDifyToolsBuilder._provider_id(tool_config),
tool_name=tool_config.tool_name,
tool_parameters=dict(tool_config.runtime_parameters),
@@ -350,16 +343,24 @@ class WorkflowAgentDifyToolsBuilder:
@staticmethod
def _provider_key(tool_config: AgentSoulDifyToolConfig) -> tuple[ToolProviderType, str]:
return (tool_config.provider_type, WorkflowAgentDifyToolsBuilder._provider_id(tool_config))
return (
ToolProviderType.value_of(tool_config.provider_type),
WorkflowAgentDifyToolsBuilder._provider_id(tool_config),
)
@staticmethod
def _tool_layer_destination(tool_config: AgentSoulDifyToolConfig) -> Literal["plugin", "core"]:
provider_type = tool_config.provider_type
provider_type = ToolProviderType.value_of(tool_config.provider_type)
if provider_type is ToolProviderType.PLUGIN or (
provider_type is ToolProviderType.BUILT_IN and _is_plugin_provider_id(tool_config.provider_id)
):
return "plugin"
if provider_type in _CORE_TOOL_PROVIDER_TYPES:
if provider_type in {
ToolProviderType.BUILT_IN,
ToolProviderType.API,
ToolProviderType.WORKFLOW,
ToolProviderType.MCP,
}:
return "core"
if provider_type is ToolProviderType.DATASET_RETRIEVAL:
raise WorkflowAgentDifyToolsBuildError(
@@ -416,7 +417,7 @@ class WorkflowAgentDifyToolsBuilder:
) -> DifyCoreToolConfig:
parameters = self._prepared_parameters(tool_runtime)
return DifyCoreToolConfig(
provider_type=_CORE_TOOL_PROVIDER_TYPES[tool_config.provider_type],
provider_type=cast(DifyCoreToolProviderType, tool_config.provider_type),
provider_id=self._provider_id(tool_config),
tool_name=tool_config.tool_name or exposed_name,
credential_id=tool_config.credential_ref.id if tool_config.credential_ref else None,
@@ -250,6 +250,7 @@ class WorkflowAgentRuntimeRequestBuilder:
agent_config_version_kind="snapshot",
agent_mode=self._agent_backend_agent_mode(context.dify_context.invoke_from),
invoke_from=cast(DifyExecutionContextInvokeFrom, context.dify_context.invoke_from.value),
trace_id=context.dify_context.trace_session_id,
),
agent_soul_prompt=soul_prompt or None,
workflow_node_job_prompt=workflow_job_prompt,
+346
View File
@@ -0,0 +1,346 @@
"""Enforce focused KnowledgeFS coverage without hiding integration-critical glue.
The primary threshold aggregates statement and branch coverage for every Dify
module owned by the KnowledgeFS integration. Large pre-existing Dify modules
that only contain narrow integration hooks are checked with changed-line
coverage instead, so unrelated legacy code cannot dilute or inflate the gate.
"""
from __future__ import annotations
import argparse
import json
import logging
import os
import re
import subprocess
from dataclasses import dataclass
from pathlib import Path
from typing import TypedDict, cast
WORKSPACE_ROOT = Path(__file__).resolve().parents[2]
NON_CORE_COVERAGE_ALLOWLIST = frozenset(
{
"api/dev/check_knowledge_fs_coverage.py",
"api/dev/generate_knowledge_fs_contract.py",
"api/dev/knowledge_fs_product_contract.py",
"api/migrations/versions/2026_07_21_1200-a4e7c2f91b30_add_knowledge_fs_control_plane.py",
"api/migrations/versions/2026_07_21_1300-b7f2a9d41c60_add_knowledge_fs_cutover.py",
"api/migrations/versions/2026_07_21_1400-c8e31b7d52a4_add_knowledge_fs_cleanup_authorization.py",
}
)
HUNK_HEADER = re.compile(r"^@@ -\d+(?:,\d+)? \+(\d+)(?:,\d+)? @@")
logger = logging.getLogger(__name__)
class CoverageSummary(TypedDict):
"""Coverage.py counts required by the aggregate gate."""
covered_lines: int
num_statements: int
covered_branches: int
num_branches: int
class CoverageFile(TypedDict):
"""Per-file coverage data emitted by ``coverage json``."""
executed_lines: list[int]
missing_lines: list[int]
summary: CoverageSummary
class CoverageReport(TypedDict):
"""Relevant top-level shape of a coverage.py JSON report."""
files: dict[str, CoverageFile]
@dataclass(frozen=True, slots=True)
class CoverageTotals:
"""Covered and measurable units for one gate surface."""
covered: int
total: int
@property
def percent(self) -> float:
return 100.0 if self.total == 0 else self.covered * 100 / self.total
class CoverageGateError(RuntimeError):
"""Raised when coverage input is incomplete or below its threshold."""
def main() -> None:
"""Validate focused module coverage and changed integration glue."""
logging.basicConfig(level=logging.INFO, format="%(message)s")
parser = argparse.ArgumentParser()
parser.add_argument("--coverage-json", type=Path, required=True)
parser.add_argument("--glue-manifest", type=Path, required=True)
parser.add_argument("--workspace-root", type=Path, default=WORKSPACE_ROOT)
parser.add_argument("--base", default="")
parser.add_argument("--minimum", type=float, default=90.0)
parser.add_argument("--glue-minimum", type=float, default=90.0)
args = parser.parse_args()
workspace_root = args.workspace_root.resolve()
report = load_coverage_report(args.coverage_json)
core_totals = validate_core_coverage(report, workspace_root=workspace_root, minimum=args.minimum)
base = resolve_diff_base(workspace_root, args.base)
glue_paths = load_glue_coverage_paths(args.glue_manifest, workspace_root=workspace_root)
changed_lines = collect_changed_glue_lines(workspace_root, base, glue_paths=glue_paths)
glue_totals = validate_changed_glue_coverage(
report,
changed_lines=changed_lines,
minimum=args.glue_minimum,
)
logger.info(
"Dify KnowledgeFS coverage passed: core lines+branches %.2f%% (%d/%d); changed glue lines %.2f%% (%d/%d)",
core_totals.percent,
core_totals.covered,
core_totals.total,
glue_totals.percent,
glue_totals.covered,
glue_totals.total,
)
def load_coverage_report(path: Path) -> CoverageReport:
"""Load the detailed JSON report used by both coverage checks."""
if not path.is_file():
raise CoverageGateError(f"coverage JSON does not exist: {path}")
document = json.loads(path.read_text())
if not isinstance(document, dict) or not isinstance(document.get("files"), dict):
raise CoverageGateError(f"coverage JSON has no files object: {path}")
return cast(CoverageReport, document)
def load_glue_coverage_paths(path: Path, *, workspace_root: Path) -> tuple[str, ...]:
"""Load the workflow's authoritative NUL-delimited integration touchpoints."""
if not path.is_file():
raise CoverageGateError(f"KnowledgeFS glue manifest does not exist: {path}")
try:
paths = tuple(item.decode() for item in path.read_bytes().split(b"\0") if item)
except UnicodeDecodeError as error:
raise CoverageGateError(f"KnowledgeFS glue manifest is not UTF-8: {path}") from error
if not paths:
raise CoverageGateError("Dify KnowledgeFS glue coverage target set is empty")
if len(paths) != len(set(paths)):
raise CoverageGateError("Dify KnowledgeFS glue coverage manifest contains duplicate paths")
invalid_paths = [
candidate
for candidate in paths
if not candidate.startswith("api/")
or not candidate.endswith(".py")
or not (workspace_root / candidate).is_file()
]
if invalid_paths:
raise CoverageGateError(f"Dify KnowledgeFS glue coverage paths are invalid: {', '.join(invalid_paths)}")
return paths
def is_core_coverage_path(path: str) -> bool:
"""Return whether a repository-relative path belongs to the focused aggregate."""
if not path.endswith(".py"):
return False
if path in {
"api/commands/knowledge_fs.py",
"api/configs/extra/knowledge_fs_config.py",
"api/extensions/ext_knowledge_fs_observability.py",
"api/services/knowledge_fs_capability.py",
}:
return True
if path.startswith(
(
"api/controllers/console/knowledge_fs/",
"api/controllers/service_api/knowledge_fs/",
"api/core/tools/builtin_tool/providers/knowledge_fs/",
"api/services/knowledge_fs/",
)
):
return True
filename = path.rsplit("/", maxsplit=1)[-1]
return (
path.startswith("api/models/")
and filename.startswith("knowledge_fs")
or path.startswith("api/repositories/")
and "knowledge_fs" in filename
or path.startswith("api/tasks/")
and "knowledge_fs" in filename
)
def discover_core_coverage_paths(workspace_root: Path) -> tuple[str, ...]:
"""Classify every KnowledgeFS-named production file or fail closed."""
named_paths = discover_knowledge_fs_production_paths(workspace_root)
core_paths = {path for path in named_paths if is_core_coverage_path(path)}
unclassified_paths = set(named_paths) - core_paths - NON_CORE_COVERAGE_ALLOWLIST
if unclassified_paths:
raise CoverageGateError(
"unclassified Dify KnowledgeFS production files must join the core coverage scope or explicit allowlist: "
+ ", ".join(sorted(unclassified_paths))
)
if not core_paths:
raise CoverageGateError("Dify KnowledgeFS core coverage target set is empty")
return tuple(sorted(core_paths))
def discover_knowledge_fs_production_paths(workspace_root: Path) -> tuple[str, ...]:
"""Mirror the workflow's dynamic KnowledgeFS filename discovery."""
api_root = workspace_root / "api"
if not api_root.is_dir():
raise CoverageGateError(f"Dify API directory does not exist: {api_root}")
paths: set[str] = set()
for directory, child_directories, filenames in os.walk(api_root):
current_directory = Path(directory)
if current_directory == api_root:
child_directories[:] = [name for name in child_directories if name not in {".venv", "storage", "tests"}]
child_directories[:] = [name for name in child_directories if name != "__pycache__"]
for filename in filenames:
path = (current_directory / filename).relative_to(workspace_root).as_posix()
if filename.endswith(".py") and "knowledge_fs" in path:
paths.add(path)
if not paths:
raise CoverageGateError("Dify KnowledgeFS production target set is empty")
return tuple(sorted(paths))
def validate_core_coverage(
report: CoverageReport,
*,
workspace_root: Path,
minimum: float,
) -> CoverageTotals:
"""Require the exact combined line-and-branch percentage for all core files."""
paths = discover_core_coverage_paths(workspace_root)
missing_paths = [path for path in paths if path not in report["files"]]
if missing_paths:
raise CoverageGateError(f"coverage report is missing core files: {', '.join(missing_paths)}")
covered = 0
total = 0
for path in paths:
summary = report["files"][path]["summary"]
covered += summary["covered_lines"] + summary["covered_branches"]
total += summary["num_statements"] + summary["num_branches"]
if total == 0:
raise CoverageGateError("Dify KnowledgeFS core coverage has no measurable statements or branches")
totals = CoverageTotals(covered=covered, total=total)
_require_minimum(totals, minimum=minimum, label="Dify KnowledgeFS core line-and-branch coverage")
return totals
def resolve_diff_base(workspace_root: Path, preferred: str) -> str:
"""Resolve an explicit event base, falling back to the previous commit for manual runs."""
base = preferred.strip()
if not base or set(base) == {"0"}:
base = "HEAD^"
result = subprocess.run(
["git", "cat-file", "-e", f"{base}^{{commit}}"],
cwd=workspace_root,
check=False,
capture_output=True,
text=True,
)
if result.returncode != 0:
detail = result.stderr.strip() or "commit is unavailable"
raise CoverageGateError(f"cannot resolve coverage diff base {base}: {detail}")
return base
def collect_changed_glue_lines(
workspace_root: Path,
base: str,
*,
glue_paths: tuple[str, ...],
) -> dict[str, set[int]]:
"""Return added line numbers in the narrow Dify modules touched by this integration."""
result = subprocess.run(
[
"git",
"diff",
"--no-ext-diff",
"--no-color",
"--unified=0",
base,
"--",
*glue_paths,
],
cwd=workspace_root,
check=False,
capture_output=True,
text=True,
)
if result.returncode != 0:
detail = result.stderr.strip() or "git diff failed"
raise CoverageGateError(f"cannot collect KnowledgeFS glue diff from {base}: {detail}")
return parse_added_lines(result.stdout)
def parse_added_lines(diff: str) -> dict[str, set[int]]:
"""Parse repository paths and added-side line numbers from a zero-context Git diff."""
changed_lines: dict[str, set[int]] = {}
current_path: str | None = None
current_line: int | None = None
for raw_line in diff.splitlines():
if raw_line.startswith("diff --git "):
current_line = None
continue
if raw_line.startswith("+++ "):
candidate = raw_line[4:]
current_path = candidate[2:] if candidate.startswith("b/") else None
if current_path is not None:
changed_lines.setdefault(current_path, set())
current_line = None
continue
if raw_line.startswith("@@ "):
match = HUNK_HEADER.match(raw_line)
current_line = int(match.group(1)) if match is not None else None
continue
if current_path is None or current_line is None:
continue
if raw_line.startswith("+"):
changed_lines[current_path].add(current_line)
current_line += 1
elif raw_line.startswith("-") or raw_line.startswith("\\"):
continue
else:
current_line += 1
return changed_lines
def validate_changed_glue_coverage(
report: CoverageReport,
*,
changed_lines: dict[str, set[int]],
minimum: float,
) -> CoverageTotals:
"""Require added executable glue lines to be exercised by the focused unit suite."""
covered = 0
total = 0
for path, lines in sorted(changed_lines.items()):
if not lines:
continue
file_coverage = report["files"].get(path)
if file_coverage is None:
raise CoverageGateError(f"coverage report is missing changed glue file: {path}")
executed_lines = set(file_coverage["executed_lines"])
executable_lines = executed_lines | set(file_coverage["missing_lines"])
changed_executable_lines = lines & executable_lines
covered += len(changed_executable_lines & executed_lines)
total += len(changed_executable_lines)
totals = CoverageTotals(covered=covered, total=total)
_require_minimum(totals, minimum=minimum, label="Dify KnowledgeFS changed-glue line coverage")
return totals
def _require_minimum(totals: CoverageTotals, *, minimum: float, label: str) -> None:
if totals.percent + 1e-12 < minimum:
raise CoverageGateError(f"{label} is {totals.percent:.2f}%; minimum {minimum:.2f}%")
if __name__ == "__main__":
main()
+534 -187
View File
@@ -1,7 +1,14 @@
"""Validate Dify Console KnowledgeFS declarations against a pinned OpenAPI document.
"""Pin the in-repository KnowledgeFS contract and validate every Dify product operation.
The OpenAPI document is exported only during explicit development validation. Runtime declarations live with Dify
product policy; this module validates their transport metadata without generating a complete operation catalog.
The lock is intentionally independent of the enclosing Dify commit: it records the staged ``knowledge-fs/`` tree,
the complete generated OpenAPI document, both explicit product-operation manifests, and the active Capability v2
profile and deterministic public-key vector. The full OpenAPI hash
covers request/response schemas, status codes, security, deprecation, and stream metadata. Field-level validation
cross-checks the Dify product registry, Python Capability issuer, TypeScript request guard, and exported OpenAPI;
each product operation must be ready or an explicit gap, and KFS-only activation remains explicitly internal.
Contract export reads the working tree only after proving it matches the staged KnowledgeFS index. This keeps the
OpenAPI bytes and auth manifest aligned with the exact subtree tree ID that will be reviewed and committed.
"""
from __future__ import annotations
@@ -12,29 +19,37 @@ import json
import subprocess
import sys
import tempfile
from copy import deepcopy
from pathlib import Path
from typing import Any, Literal, TypedDict
from typing import Any, Literal, TypedDict, cast
import jwt
from cryptography.hazmat.primitives.asymmetric.rsa import RSAPublicKey
from jwt.algorithms import RSAAlgorithm
API_ROOT = Path(__file__).resolve().parents[1]
if str(API_ROOT) not in sys.path:
sys.path.insert(0, str(API_ROOT))
from dev.knowledge_fs_product_contract import (
capability_operation_runtime_contracts,
parse_capability_operation_policy,
parse_product_operation_gap_manifest,
parse_product_operation_manifest,
product_operation_runtime_contracts,
validate_product_operation_contracts,
)
WORKSPACE_ROOT = API_ROOT.parent
LOCK_PATH = API_ROOT / "knowledge-fs-contract.lock.json"
DEFAULT_REPOSITORY = WORKSPACE_ROOT.parent / "knowledge-fs"
KNOWLEDGE_FS_DIRECTORY = "knowledge-fs"
CAPABILITY_V2_AUTH_MANIFEST_RELATIVE_PATH = Path("contracts/dify-capability-v2-auth-profile.json")
CAPABILITY_V2_AUTH_TEST_VECTOR_RELATIVE_PATH = Path("contracts/dify-capability-v2-test-vector.json")
UPSTREAM_PROVENANCE_RELATIVE_PATH = Path("upstream-provenance.json")
LOCK_RELATIVE_PATH = Path("api/knowledge-fs-contract.lock.json")
PRODUCT_OPERATIONS_RELATIVE_PATH = Path("api/knowledge-fs-product-operations.json")
PRODUCT_OPERATION_GAPS_RELATIVE_PATH = Path("api/knowledge-fs-product-operation-gaps.json")
OPENAPI_METHODS = ("delete", "get", "head", "options", "patch", "post", "put", "trace")
PROXY_METHODS = frozenset({"delete", "get", "patch", "post", "put"})
CONSOLE_PROXY_ERROR_SCHEMA_NAME = "ConsoleProxyError"
CONSOLE_PROXY_ERROR_SCHEMA: dict[str, Any] = {
"type": "object",
"required": ["code", "message", "status"],
"properties": {
"code": {"type": "string"},
"message": {"type": "string"},
"status": {"type": "integer"},
},
}
LOCK_SCHEMA_VERSION = 5
class ContractDeclaration(TypedDict):
@@ -49,7 +64,18 @@ class ContractDeclaration(TypedDict):
request_headers: tuple[str, ...]
response_headers: tuple[str, ...]
response_media_types: tuple[str, ...]
error_status_map: tuple[tuple[int, int], ...]
class ContractLock(TypedDict):
"""Content-addressed contract inputs that must move together."""
schemaVersion: int
subtreeTree: str
openapiSha256: str
capabilityV2AuthManifestSha256: str
capabilityV2AuthTestVectorSha256: str
productOperationManifestSha256: str
productOperationGapManifestSha256: str
type DeclarationField = Literal[
@@ -76,68 +102,507 @@ DECLARATION_FIELDS: tuple[DeclarationField, ...] = (
def main() -> None:
"""Update or verify the pin and validate Console declarations against its OpenAPI document."""
"""Update or verify the monorepo pin and validate Dify product declarations."""
parser = argparse.ArgumentParser()
mode = parser.add_mutually_exclusive_group()
mode = parser.add_mutually_exclusive_group(required=True)
mode.add_argument("--check", action="store_true")
mode.add_argument("--update-lock", action="store_true")
parser.add_argument("--repository", type=Path, default=DEFAULT_REPOSITORY)
parser.add_argument("--output-openapi", type=Path)
parser.add_argument("--workspace-root", type=Path, default=WORKSPACE_ROOT)
args = parser.parse_args()
repository = args.repository.resolve()
lock = json.loads(LOCK_PATH.read_text())
tracked_changes = run("git", "status", "--porcelain", "--untracked-files=no", cwd=repository).strip()
if tracked_changes:
raise RuntimeError("KnowledgeFS checkout must not contain tracked changes during contract export")
workspace_root = args.workspace_root.resolve()
knowledge_fs_root = workspace_root / KNOWLEDGE_FS_DIRECTORY
lock_path = workspace_root / LOCK_RELATIVE_PATH
capability_v2_auth_manifest_path = knowledge_fs_root / CAPABILITY_V2_AUTH_MANIFEST_RELATIVE_PATH
capability_v2_auth_test_vector_path = knowledge_fs_root / CAPABILITY_V2_AUTH_TEST_VECTOR_RELATIVE_PATH
product_operations_path = workspace_root / PRODUCT_OPERATIONS_RELATIVE_PATH
product_operation_gaps_path = workspace_root / PRODUCT_OPERATION_GAPS_RELATIVE_PATH
upstream_provenance_path = knowledge_fs_root / UPSTREAM_PROVENANCE_RELATIVE_PATH
ensure_clean_knowledge_fs_worktree(workspace_root)
ensure_contract_inputs_exist(
knowledge_fs_root=knowledge_fs_root,
capability_v2_auth_manifest_path=capability_v2_auth_manifest_path,
capability_v2_auth_test_vector_path=capability_v2_auth_test_vector_path,
product_operations_path=product_operations_path,
product_operation_gaps_path=product_operation_gaps_path,
upstream_provenance_path=upstream_provenance_path,
)
commit = run("git", "rev-parse", "HEAD", cwd=repository).strip()
if not args.update_lock and commit != lock["commit"]:
raise RuntimeError(
f"KnowledgeFS checkout mismatch: expected {lock['commit']}, received {commit}. "
"Use the pinned commit or pass --update-lock intentionally."
)
subtree_tree = staged_subtree_tree(workspace_root)
capability_v2_auth_manifest_content = capability_v2_auth_manifest_path.read_bytes()
capability_v2_auth_test_vector_content = capability_v2_auth_test_vector_path.read_bytes()
product_operation_manifest_content = product_operations_path.read_bytes()
product_operation_gap_manifest_content = product_operation_gaps_path.read_bytes()
capability_v2_auth_manifest = load_json_object(capability_v2_auth_manifest_path)
validate_capability_v2_auth_manifest(capability_v2_auth_manifest)
validate_capability_v2_auth_test_vector(
load_json_object(capability_v2_auth_test_vector_path),
capability_v2_auth_manifest,
)
validate_upstream_provenance(load_json_object(upstream_provenance_path))
product_manifest = parse_product_operation_manifest(load_json_object(product_operations_path))
product_gap_manifest = parse_product_operation_gap_manifest(load_json_object(product_operation_gaps_path))
product_runtime_operations = product_operation_runtime_contracts()
capability_runtime_operations = capability_operation_runtime_contracts()
declarations = console_contract_declarations()
with tempfile.TemporaryDirectory(prefix="dify-knowledge-fs-contract-") as directory:
openapi_path = Path(directory) / "knowledge-fs.openapi.json"
capability_policy_path = Path(directory) / "dify-capability-v2-operations.json"
subprocess.run(
["pnpm", "openapi:export", "--", "--output", str(openapi_path)],
cwd=repository,
cwd=knowledge_fs_root,
check=True,
)
subprocess.run(
["pnpm", "capability:export", "--", "--output", str(capability_policy_path)],
cwd=knowledge_fs_root,
check=True,
)
openapi_content = openapi_path.read_bytes()
openapi_sha256 = sha256(openapi_content)
if not args.update_lock and openapi_sha256 != lock["openapiSha256"]:
raise RuntimeError(
f"KnowledgeFS OpenAPI hash mismatch: expected {lock['openapiSha256']}, received {openapi_sha256}"
)
capability_policy = parse_capability_operation_policy(load_json_object(capability_policy_path))
document: dict[str, Any] = json.loads(openapi_content)
declarations = console_contract_declarations()
validate_product_operation_contracts(
capability_operations=capability_runtime_operations,
capability_policy=capability_policy,
document=document,
gap_manifest=product_gap_manifest,
manifest=product_manifest,
product_operations=product_runtime_operations,
)
validate_declarations(document, declarations)
if args.output_openapi:
filtered_document = filter_openapi_document(document, declarations)
filtered_document["x-dify-source-openapi-sha256"] = openapi_sha256
filtered_document["x-dify-console-declarations-sha256"] = contract_declarations_sha256(declarations)
args.output_openapi.parent.mkdir(parents=True, exist_ok=True)
args.output_openapi.write_text(json.dumps(filtered_document, indent=2) + "\n")
expected_lock: ContractLock = {
"schemaVersion": LOCK_SCHEMA_VERSION,
"subtreeTree": subtree_tree,
"openapiSha256": sha256(openapi_content),
"capabilityV2AuthManifestSha256": sha256(capability_v2_auth_manifest_content),
"capabilityV2AuthTestVectorSha256": sha256(capability_v2_auth_test_vector_content),
"productOperationManifestSha256": sha256(product_operation_manifest_content),
"productOperationGapManifestSha256": sha256(product_operation_gap_manifest_content),
}
if args.update_lock:
LOCK_PATH.write_text(
json.dumps(
{
"commit": commit,
"openapiSha256": openapi_sha256,
"repository": lock["repository"],
},
indent=2,
lock_path.write_text(json.dumps(expected_lock, indent=2) + "\n")
return
received_lock = parse_contract_lock(load_json_object(lock_path))
lock_fields = (
(
"capabilityV2AuthManifestSha256",
received_lock["capabilityV2AuthManifestSha256"],
expected_lock["capabilityV2AuthManifestSha256"],
),
(
"capabilityV2AuthTestVectorSha256",
received_lock["capabilityV2AuthTestVectorSha256"],
expected_lock["capabilityV2AuthTestVectorSha256"],
),
(
"productOperationManifestSha256",
received_lock["productOperationManifestSha256"],
expected_lock["productOperationManifestSha256"],
),
(
"productOperationGapManifestSha256",
received_lock["productOperationGapManifestSha256"],
expected_lock["productOperationGapManifestSha256"],
),
("schemaVersion", received_lock["schemaVersion"], expected_lock["schemaVersion"]),
("subtreeTree", received_lock["subtreeTree"], expected_lock["subtreeTree"]),
("openapiSha256", received_lock["openapiSha256"], expected_lock["openapiSha256"]),
)
for field, received_value, expected_value in lock_fields:
if received_value != expected_value:
raise RuntimeError(
f"KnowledgeFS contract lock field {field} drifted: "
f"expected {expected_value!r}, received {received_value!r}. "
"Run --update-lock intentionally after reviewing the staged subtree and contract changes."
)
+ "\n"
def ensure_contract_inputs_exist(
*,
knowledge_fs_root: Path,
capability_v2_auth_manifest_path: Path,
capability_v2_auth_test_vector_path: Path,
product_operations_path: Path,
product_operation_gaps_path: Path,
upstream_provenance_path: Path,
) -> None:
"""Fail with a stable error before invoking package tooling when a contract input is absent."""
required_paths = (
knowledge_fs_root / "package.json",
capability_v2_auth_manifest_path,
capability_v2_auth_test_vector_path,
product_operations_path,
product_operation_gaps_path,
upstream_provenance_path,
)
missing_paths = [path for path in required_paths if not path.is_file()]
if missing_paths:
missing = ", ".join(str(path) for path in missing_paths)
raise RuntimeError(f"KnowledgeFS contract input is missing: {missing}")
def ensure_clean_knowledge_fs_worktree(workspace_root: Path) -> None:
"""Require exported KnowledgeFS files to exactly match the staged index tree.
Staged changes are expected during intentional lock updates. Unstaged tracked changes and untracked files are
rejected because the export process reads the working tree while the tree ID is calculated from the index.
"""
subtree_path = f"{KNOWLEDGE_FS_DIRECTORY}/"
unstaged = subprocess.run(
["git", "diff", "--quiet", "--", subtree_path],
cwd=workspace_root,
check=False,
)
if unstaged.returncode > 1:
raise RuntimeError("git diff failed while validating the staged KnowledgeFS subtree")
untracked = run(
"git",
"ls-files",
"--others",
"--exclude-standard",
"--",
subtree_path,
cwd=workspace_root,
).strip()
if unstaged.returncode != 0 or untracked:
raise RuntimeError(
"knowledge-fs/ contains unstaged or untracked changes; stage or remove them before contract export"
)
def staged_subtree_tree(workspace_root: Path) -> str:
"""Return the Git tree object for the staged ``knowledge-fs/`` subtree."""
return run("git", "write-tree", f"--prefix={KNOWLEDGE_FS_DIRECTORY}/", cwd=workspace_root).strip()
def load_json_object(path: Path) -> dict[str, Any]:
"""Load a JSON object and reject arrays/scalars at contract boundaries."""
value = json.loads(path.read_text())
if not isinstance(value, dict):
raise ValueError(f"KnowledgeFS contract file must contain a JSON object: {path}")
return cast(dict[str, Any], value)
def parse_contract_lock(value: dict[str, Any]) -> ContractLock:
"""Validate the compact, non-self-referential contract lock schema."""
expected_fields = {
"capabilityV2AuthManifestSha256",
"capabilityV2AuthTestVectorSha256",
"openapiSha256",
"productOperationGapManifestSha256",
"productOperationManifestSha256",
"schemaVersion",
"subtreeTree",
}
if set(value) != expected_fields:
raise ValueError(f"KnowledgeFS contract lock fields must be exactly {sorted(expected_fields)}")
if value.get("schemaVersion") != LOCK_SCHEMA_VERSION:
raise ValueError(f"KnowledgeFS contract lock schemaVersion must be {LOCK_SCHEMA_VERSION}")
for field in (
"capabilityV2AuthManifestSha256",
"capabilityV2AuthTestVectorSha256",
"openapiSha256",
"productOperationGapManifestSha256",
"productOperationManifestSha256",
"subtreeTree",
):
field_value = value.get(field)
expected_length = 40 if field == "subtreeTree" else 64
if (
not isinstance(field_value, str)
or len(field_value) != expected_length
or any(character not in "0123456789abcdef" for character in field_value)
):
raise ValueError(f"KnowledgeFS contract lock field {field} has an invalid digest")
return cast(ContractLock, value)
def validate_required_product_operations(document: dict[str, Any], required_operation_ids: list[str]) -> None:
"""Fail when the pinned KFS OpenAPI omits an operation required by the Dify product."""
available_operation_ids = [
operation_id
for path_item in document.get("paths", {}).values()
for method in OPENAPI_METHODS
if isinstance(path_item, dict)
for operation in (path_item.get(method),)
if isinstance(operation, dict)
for operation_id in (operation.get("operationId"),)
if isinstance(operation_id, str) and operation_id
]
for operation_id in required_operation_ids:
count = available_operation_ids.count(operation_id)
if count != 1:
raise ValueError(
f"KnowledgeFS OpenAPI required product operation {operation_id} must occur exactly once; found {count}"
)
def validate_capability_v2_auth_manifest(value: dict[str, Any]) -> None:
"""Validate the active production RS256 profile consumed by both Dify and KnowledgeFS."""
expected_fields = {
"active",
"audience",
"callerProfiles",
"claimBindings",
"issuer",
"lifecycle",
"maxTtlSeconds",
"productionReady",
"profileId",
"protectedHeader",
"requiredClaims",
"resourceContract",
"runtimeAssembly",
"schemaVersion",
"signatureAlgorithms",
"tokenKind",
}
if set(value) != expected_fields:
raise ValueError(f"KnowledgeFS Capability v2 auth manifest fields must be exactly {sorted(expected_fields)}")
fixed_values = {
"active": True,
"audience": "knowledge-fs",
"issuer": "dify-control-plane",
"lifecycle": "active",
"maxTtlSeconds": 60,
"productionReady": True,
"profileId": "dify-capability-v2",
"schemaVersion": 3,
"signatureAlgorithms": ["RS256"],
"tokenKind": "jwt",
}
for field, expected in fixed_values.items():
if value.get(field) != expected:
raise ValueError(f"KnowledgeFS Capability v2 auth manifest field {field} must be {expected!r}")
if value.get("protectedHeader") != {
"algorithm": "RS256",
"keyIdClaim": "kid",
"keyIdRequired": True,
"type": "JWT",
}:
raise ValueError("KnowledgeFS Capability v2 protected header contract is invalid")
required_claims = [
"action",
"actor",
"aud",
"authz_revision",
"azp",
"caller_kind",
"cap_ver",
"content_policy_revision",
"content_scope_ids",
"control_space_id",
"exp",
"grant_id",
"iat",
"iss",
"jti",
"namespace_id",
"nbf",
"resource",
"sub",
"trace_id",
]
if value.get("requiredClaims") != required_claims:
raise ValueError("KnowledgeFS Capability v2 required claims are invalid")
if value.get("claimBindings") != {
"action": "action",
"callerKind": "caller_kind",
"controlSpace": "control_space_id",
"namespace": "namespace_id",
"resource": "resource",
"resourceParent": "resource.parent_id",
"subject": "sub",
}:
raise ValueError("KnowledgeFS Capability v2 claim bindings are invalid")
if value.get("resourceContract") != {
"fields": ["id", "parent_id", "type"],
"parentForbiddenFor": ["namespace", "knowledge_space"],
"parentRequiredFor": ["document", "job", "query", "research_task", "source", "upload_session"],
}:
raise ValueError("KnowledgeFS Capability v2 resource contract is invalid")
if value.get("callerProfiles") != {
"agent": {"authorizedParty": "dify-agent", "subjectPrefix": "dify-app:"},
"interactive": {"authorizedParty": "dify-console", "subjectPrefix": "dify-account:"},
"internal_worker": {"authorizedParty": "dify-worker", "subjectPrefix": "dify-worker:"},
"mcp": {"authorizedParty": "dify-mcp", "subjectPrefix": "dify-mcp-session:"},
"service": {"authorizedParty": "dify-service-api", "subjectPrefix": "dify-kfs-credential:"},
"workflow": {"authorizedParty": "dify-workflow", "subjectPrefix": "dify-app:"},
}:
raise ValueError("KnowledgeFS Capability v2 caller profiles are invalid")
if value.get("runtimeAssembly") != {
"failClosed": True,
"keySelection": "kid",
"maximumPublishedKeys": 3,
"verificationKeySource": "jwks",
}:
raise ValueError("KnowledgeFS Capability v2 runtime assembly is invalid")
def validate_capability_v2_auth_test_vector(
value: dict[str, Any],
manifest: dict[str, Any],
) -> None:
"""Verify the deterministic public-key vector and every security-sensitive binding."""
expected_fields = {
"algorithm",
"audience",
"expectedClaims",
"expectedPrincipal",
"issuer",
"operation",
"profileId",
"protectedHeader",
"publicJwk",
"schemaVersion",
"testOnly",
"token",
"ttlSeconds",
}
if set(value) != expected_fields:
raise ValueError(f"KnowledgeFS Capability v2 test vector fields must be exactly {sorted(expected_fields)}")
if (
value.get("schemaVersion") != 2
or value.get("profileId") != manifest.get("profileId")
or value.get("testOnly") is not True
or value.get("algorithm") != "RS256"
or value.get("issuer") != manifest.get("issuer")
or value.get("audience") != manifest.get("audience")
or value.get("ttlSeconds") != manifest.get("maxTtlSeconds")
):
raise ValueError("KnowledgeFS Capability v2 test vector does not match the active profile")
protected_header = value.get("protectedHeader")
if not isinstance(protected_header, dict) or set(protected_header) != {"alg", "kid", "typ"}:
raise ValueError("KnowledgeFS Capability v2 test vector protected header is invalid")
kid = protected_header.get("kid")
if protected_header.get("alg") != "RS256" or protected_header.get("typ") != "JWT" or not _is_non_blank(kid):
raise ValueError("KnowledgeFS Capability v2 test vector protected header is invalid")
public_jwk = value.get("publicJwk")
if not isinstance(public_jwk, dict) or set(public_jwk) != {"alg", "e", "kid", "kty", "n", "use"}:
raise ValueError("KnowledgeFS Capability v2 test vector public JWK is invalid")
if (
public_jwk.get("alg") != "RS256"
or public_jwk.get("kid") != kid
or public_jwk.get("kty") != "RSA"
or public_jwk.get("use") != "sig"
or not _is_non_blank(public_jwk.get("e"))
or not _is_non_blank(public_jwk.get("n"))
):
raise ValueError("KnowledgeFS Capability v2 test vector public JWK is invalid")
claims = value.get("expectedClaims")
required_claims = manifest.get("requiredClaims")
if not isinstance(claims, dict) or not isinstance(required_claims, list) or set(claims) != set(required_claims):
raise ValueError("KnowledgeFS Capability v2 test vector claims do not match the active profile")
operation = value.get("operation")
if not isinstance(operation, dict) or set(operation) != {"action", "method", "operationId", "requestPath"}:
raise ValueError("KnowledgeFS Capability v2 test vector operation is invalid")
resource = claims.get("resource")
if not isinstance(resource, dict) or set(resource) != {"id", "parent_id", "type"}:
raise ValueError("KnowledgeFS Capability v2 test vector resource is invalid")
expected_operation = {
"action": "documents.read",
"method": "GET",
"operationId": "getDocument",
"requestPath": "/knowledge-spaces/space-contract-vector/documents/document-contract-vector",
}
if operation != expected_operation:
raise ValueError("KnowledgeFS Capability v2 test vector operation binding is invalid")
exact_claims = {
"action": operation["action"],
"aud": value["audience"],
"caller_kind": "interactive",
"cap_ver": 2,
"control_space_id": "control-space-contract-vector",
"iss": value["issuer"],
"namespace_id": "workspace-contract-vector",
"resource": {
"id": "document-contract-vector",
"parent_id": "space-contract-vector",
"type": "document",
},
"sub": "dify-account:account-contract-vector",
}
for field, expected in exact_claims.items():
if claims.get(field) != expected:
raise ValueError(f"KnowledgeFS Capability v2 test vector claim {field} is invalid")
if claims.get("actor") != claims["sub"] or claims.get("azp") != "dify-console":
raise ValueError("KnowledgeFS Capability v2 test vector caller binding is invalid")
issued_at = claims.get("iat")
not_before = claims.get("nbf")
expires_at = claims.get("exp")
if (
not isinstance(issued_at, int)
or isinstance(issued_at, bool)
or not_before != issued_at
or not isinstance(expires_at, int)
or isinstance(expires_at, bool)
or expires_at - issued_at != value["ttlSeconds"]
):
raise ValueError("KnowledgeFS Capability v2 test vector TTL is invalid")
expected_principal = {
"callerKind": claims["caller_kind"],
"subject": {
"scopes": ["knowledge-spaces:read"],
"subjectId": claims["sub"],
"tenantId": claims["namespace_id"],
},
}
if value.get("expectedPrincipal") != expected_principal:
raise ValueError("KnowledgeFS Capability v2 test vector principal is invalid")
token = value.get("token")
if not isinstance(token, str) or not _is_non_blank(token):
raise ValueError("KnowledgeFS Capability v2 test vector token is invalid")
try:
verification_key = RSAAlgorithm.from_jwk(public_jwk)
if not isinstance(verification_key, RSAPublicKey):
raise ValueError("Capability vector verification key is not RSA public material")
header = jwt.get_unverified_header(token)
decoded_claims = jwt.decode(
token,
verification_key,
algorithms=["RS256"],
audience=cast(str, value["audience"]),
issuer=cast(str, value["issuer"]),
options={"verify_exp": False, "verify_iat": False, "verify_nbf": False},
)
except (jwt.PyJWTError, TypeError, ValueError) as exc:
raise ValueError("KnowledgeFS Capability v2 test vector signature is invalid") from exc
if header != protected_header or decoded_claims != claims:
raise ValueError("KnowledgeFS Capability v2 test vector token content drifted")
def _is_non_blank(value: object) -> bool:
return isinstance(value, str) and bool(value.strip()) and value == value.strip()
def validate_upstream_provenance(value: dict[str, Any]) -> None:
"""Validate the imported-source provenance that is itself covered by the subtree tree ID."""
expected_fields = {"commit", "release", "repository", "schemaVersion"}
if set(value) != expected_fields:
raise ValueError(f"KnowledgeFS upstream provenance fields must be exactly {sorted(expected_fields)}")
if value.get("schemaVersion") != 1:
raise ValueError("KnowledgeFS upstream provenance must use schemaVersion 1")
repository = value.get("repository")
commit = value.get("commit")
if not isinstance(repository, str) or not repository.startswith("https://"):
raise ValueError("KnowledgeFS upstream provenance repository must be an HTTPS URL")
if (
not isinstance(commit, str)
or len(commit) != 40
or any(character not in "0123456789abcdef" for character in commit)
):
raise ValueError("KnowledgeFS upstream provenance commit must be a lowercase full Git SHA")
if value.get("release") is not None and not isinstance(value["release"], str):
raise ValueError("KnowledgeFS upstream provenance release must be null or a string")
def validate_declarations(document: dict[str, Any], declarations: tuple[ContractDeclaration, ...]) -> None:
"""Validate Dify Console declarations against matching pinned OpenAPI operations."""
operations_by_id: dict[str, list[tuple[str, str, dict[str, Any], dict[str, Any]]]] = {}
@@ -168,10 +633,11 @@ def validate_declarations(document: dict[str, Any], declarations: tuple[Contract
raise ValueError(f"KnowledgeFS OpenAPI path must be absolute: {path}")
if method not in PROXY_METHODS:
raise ValueError(f"KnowledgeFS proxy does not support {method.upper()} {path}")
expected: dict[DeclarationField, object] = {
expected: ContractDeclaration = {
"operation_id": operation_id,
"method": method.upper(),
"path": path[1:],
"required_scope": required_scope(operation),
"required_scope": required_scope(document, operation),
"response_kind": response_kind(operation),
"max_response_bytes": required_max_response_bytes(operation),
"request_headers": request_header_names(path_item, operation),
@@ -186,136 +652,12 @@ def validate_declarations(document: dict[str, Any], declarations: tuple[Contract
f"KnowledgeFS operation {operation_id} field {field} drifted: "
f"expected {expected_value!r}, received {received_value!r}"
)
validate_error_status_map(operation_id, declaration["error_status_map"])
def filter_openapi_document(
document: dict[str, Any],
declarations: tuple[ContractDeclaration, ...],
) -> dict[str, Any]:
"""Return a code-generation document containing only Console-allowlisted operations."""
filtered_document: dict[str, Any] = {
key: value for key, value in document.items() if key not in {"components", "paths"}
}
source_paths = document.get("paths", {})
filtered_paths: dict[str, Any] = {}
for declaration in declarations:
path = f"/{declaration['path']}"
method = declaration["method"].lower()
source_path_item = source_paths[path]
path_metadata = {key: value for key, value in source_path_item.items() if key not in OPENAPI_METHODS}
filtered_path_item = filtered_paths.setdefault(path, path_metadata)
filtered_operation = deepcopy(source_path_item[method])
_rewrite_proxy_error_responses(filtered_operation, declaration["error_status_map"])
filtered_path_item[method] = filtered_operation
filtered_document["paths"] = filtered_paths
source_components = document.get("components", {})
filtered_components = {key: value for key, value in source_components.items() if key != "schemas"}
source_schemas = source_components.get("schemas", {})
available_schemas = {**source_schemas, CONSOLE_PROXY_ERROR_SCHEMA_NAME: CONSOLE_PROXY_ERROR_SCHEMA}
schema_names = _referenced_schema_names(filtered_paths, available_schemas)
filtered_components["schemas"] = {
name: schema for name, schema in available_schemas.items() if name in schema_names
}
filtered_document["components"] = filtered_components
return filtered_document
def validate_error_status_map(operation_id: str, error_status_map: tuple[tuple[int, int], ...]) -> None:
"""Validate the status normalization advertised by one Console operation."""
upstream_statuses: set[int] = set()
for upstream_status, console_status in error_status_map:
if upstream_status in upstream_statuses:
raise ValueError(f"KnowledgeFS operation {operation_id} has duplicate error status: {upstream_status}")
if not 400 <= upstream_status <= 599 or not 400 <= console_status <= 599:
raise ValueError(f"KnowledgeFS operation {operation_id} has invalid error status mapping")
upstream_statuses.add(upstream_status)
def _rewrite_proxy_error_responses(
operation: dict[str, Any],
error_status_map: tuple[tuple[int, int], ...],
) -> None:
responses = operation.setdefault("responses", {})
proxy_error_response = {
"description": "Error normalized by the Dify Console KnowledgeFS proxy.",
"content": {
"application/json": {"schema": {"$ref": f"#/components/schemas/{CONSOLE_PROXY_ERROR_SCHEMA_NAME}"}}
},
}
for upstream_status, console_status in error_status_map:
existing_target = responses.get(str(console_status)) if upstream_status != console_status else None
responses.pop(str(upstream_status), None)
normalized_response: dict[str, Any] = deepcopy(proxy_error_response)
existing_schema = (
existing_target.get("content", {}).get("application/json", {}).get("schema")
if isinstance(existing_target, dict)
else None
)
if existing_schema is not None:
normalized_response["content"]["application/json"]["schema"] = {
"oneOf": [
deepcopy(existing_schema),
{"$ref": f"#/components/schemas/{CONSOLE_PROXY_ERROR_SCHEMA_NAME}"},
]
}
responses[str(console_status)] = normalized_response
def _referenced_schema_names(value: Any, schemas: dict[str, Any]) -> set[str]:
reference_prefix = "#/components/schemas/"
selected: set[str] = set()
pending: list[Any] = [value]
while pending:
current = pending.pop()
if isinstance(current, list):
pending.extend(current)
continue
if not isinstance(current, dict):
continue
reference = current.get("$ref")
if isinstance(reference, str) and reference.startswith(reference_prefix):
name = reference.removeprefix(reference_prefix)
if name not in selected:
if name not in schemas:
raise ValueError(f"KnowledgeFS OpenAPI references missing schema: {name}")
selected.add(name)
pending.append(schemas[name])
pending.extend(current.values())
return selected
def console_contract_declarations() -> tuple[ContractDeclaration, ...]:
"""Return transport declarations from the runtime Console operation registry."""
from services.knowledge_fs_operations import KNOWLEDGE_FS_CONSOLE_OPERATIONS
"""The P9 backend exposes only typed product controllers; the raw Console proxy is removed."""
return tuple(
{
"operation_id": operation.operation_id,
"method": operation.method,
"path": operation.path,
"required_scope": operation.required_scope,
"response_kind": operation.response_kind,
"max_response_bytes": operation.max_response_bytes,
"request_headers": operation.request_headers,
"response_headers": operation.response_headers,
"response_media_types": operation.response_media_types,
"error_status_map": operation.error_status_map,
}
for operation in KNOWLEDGE_FS_CONSOLE_OPERATIONS
)
def contract_declarations_sha256(declarations: tuple[ContractDeclaration, ...]) -> str:
"""Return a stable digest for the runtime Console operation declarations."""
content = json.dumps(declarations, separators=(",", ":"), sort_keys=True).encode()
return sha256(content)
return ()
def response_kind(operation: dict[str, Any]) -> str:
@@ -335,12 +677,17 @@ def response_media_types(operation: dict[str, Any]) -> tuple[str, ...]:
return tuple(sorted(media_types))
def required_scope(operation: dict[str, Any]) -> str | None:
def required_scope(document: dict[str, Any], operation: dict[str, Any]) -> str | None:
scope = operation.get("x-knowledge-fs-required-scope")
security = operation["security"] if "security" in operation else document.get("security")
if security == []:
if scope is not None:
raise ValueError(f"KnowledgeFS public operation must not declare a required scope: {scope}")
return None
if security != [{"bearerAuth": []}]:
raise ValueError(f"KnowledgeFS operation effective security must be exactly bearerAuth: {security}")
if scope in ("knowledge-spaces:read", "knowledge-spaces:write"):
return scope
if operation.get("security") == []:
return None
raise ValueError(f"KnowledgeFS operation has no supported required scope: {scope}")

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