Merge branch 'deploy/konwledge' of https://github.com/langgenius/dify into deploy/konwledge

This commit is contained in:
Stephen Zhou 2026-08-18 10:55:50 +08:00
commit 52cfbea82b
No known key found for this signature in database
35 changed files with 4859 additions and 43 deletions

View File

@ -1,14 +1,14 @@
from dataclasses import dataclass
from datetime import datetime
from typing import Annotated, Any
from uuid import UUID
from uuid import UUID, uuid4
from flask import request
from flask_restx import Resource
from pydantic import BaseModel, Field, field_validator, model_validator
from sqlalchemy import func, select
from sqlalchemy.orm import Session
from werkzeug.exceptions import Forbidden, NotFound
from werkzeug.exceptions import Conflict, Forbidden, NotFound
import services
from configs import dify_config
@ -32,6 +32,7 @@ from controllers.console.wraps import (
with_current_tenant_id,
with_current_user,
)
from core.db.session_factory import session_factory
from core.entities.knowledge_entities import IndexingEstimate
from core.errors.error import LLMBadRequestError, ProviderTokenNotInitError
from core.indexing_runner import IndexingRunner
@ -53,10 +54,17 @@ from models.enums import ApiTokenType, SegmentStatus
from models.provider_ids import ModelProviderID
from services.api_token_service import ApiTokenCache, get_effective_token_last_used_at
from services.app_service import AppService
from services.dataset_knowledge_fs_upgrade_service import (
KnowledgeFSUpgradeConflictError,
KnowledgeFSUpgradeNotFoundError,
KnowledgeFSUpgradeSnapshotService,
upgrade_job_response,
)
from services.dataset_ref_service import DatasetRefService
from services.dataset_service import DatasetPermissionService, DatasetService, DocumentService
from services.enterprise import rbac_service as enterprise_rbac_service
from services.enterprise.rbac_service import RBACResourceWhitelistScope, ReplaceMemberBindings
from services.knowledge_fs.product_dto import KnowledgeFSUpgradeJobResponse, KnowledgeFSUpgradeRetryResponse
from tasks.initialize_created_app_rbac_access_task import initialize_created_app_rbac_access_task
register_response_schema_models(console_ns, ApiBaseUrlResponse, SimpleResultResponse, UsageCheckResponse)
@ -356,6 +364,8 @@ register_response_schema_models(
RetrievalSettingResponse,
PartialMemberListResponse,
AutoDisableLogsResponse,
KnowledgeFSUpgradeJobResponse,
KnowledgeFSUpgradeRetryResponse,
)
@ -819,6 +829,124 @@ class DatasetApi(Resource):
raise DatasetInUseError()
@console_ns.route("/datasets/<uuid:dataset_id>/knowledge-fs-upgrades")
class DatasetKnowledgeFSUpgradeApi(Resource):
@console_ns.response(
202,
"KnowledgeFS Dataset upgrade accepted",
console_ns.models[KnowledgeFSUpgradeJobResponse.__name__],
)
@setup_required
@login_required
@account_initialization_required
@rbac_permission_required(RBACResourceScope.DATASET, RBACPermission.DATASET_EDIT)
@with_current_user
@with_current_tenant_id
def post(self, current_tenant_id: str, current_user: Account, dataset_id: UUID):
if not dify_config.KNOWLEDGE_FS_ENABLED:
raise NotFound()
dataset_id_str = str(dataset_id)
with session_factory.create_session() as session:
dataset = DatasetService.get_dataset_for_tenant(dataset_id_str, current_tenant_id, session=session)
if dataset is None:
raise NotFound("Dataset not found.")
if not dify_config.RBAC_ENABLED:
try:
DatasetService.check_dataset_permission(dataset, current_user, session)
except services.errors.account.NoPermissionError as error:
raise Forbidden(str(error)) from error
if not (current_user.has_edit_permission or current_user.is_dataset_operator):
raise Forbidden()
snapshots = KnowledgeFSUpgradeSnapshotService(session_factory.get_session_maker())
try:
job = snapshots.create(
tenant_id=current_tenant_id,
dataset_id=dataset_id_str,
requested_by_account_id=current_user.id,
idempotency_key=request.headers.get("Idempotency-Key"),
)
except KnowledgeFSUpgradeNotFoundError as error:
raise NotFound(str(error)) from error
except KnowledgeFSUpgradeConflictError as error:
raise Conflict(str(error)) from error
_enqueue_upgrade_job(snapshots, tenant_id=current_tenant_id, job_id=job.id)
return dump_response(KnowledgeFSUpgradeJobResponse, upgrade_job_response(job)), 202
@console_ns.route("/datasets/<uuid:dataset_id>/knowledge-fs-upgrades/<string:job_id>")
class DatasetKnowledgeFSUpgradeJobApi(Resource):
@console_ns.response(
200,
"KnowledgeFS Dataset upgrade status",
console_ns.models[KnowledgeFSUpgradeJobResponse.__name__],
)
@setup_required
@login_required
@account_initialization_required
@rbac_permission_required(RBACResourceScope.DATASET, RBACPermission.DATASET_READONLY)
@with_current_user
@with_current_tenant_id
def get(self, current_tenant_id: str, current_user: Account, dataset_id: UUID, job_id: str):
with session_factory.create_session() as session:
_get_accessible_dataset(dataset_id, current_tenant_id, current_user, session)
snapshots = KnowledgeFSUpgradeSnapshotService(session_factory.get_session_maker())
try:
job = snapshots.get(tenant_id=current_tenant_id, job_id=job_id)
except KnowledgeFSUpgradeNotFoundError as error:
raise NotFound(str(error)) from error
if job.old_dataset_id != str(dataset_id):
raise NotFound("Upgrade job was not found")
return dump_response(KnowledgeFSUpgradeJobResponse, upgrade_job_response(job))
@console_ns.response(
202,
"KnowledgeFS Dataset upgrade retry accepted",
console_ns.models[KnowledgeFSUpgradeRetryResponse.__name__],
)
@setup_required
@login_required
@account_initialization_required
@rbac_permission_required(RBACResourceScope.DATASET, RBACPermission.DATASET_EDIT)
@with_current_user
@with_current_tenant_id
def post(self, current_tenant_id: str, current_user: Account, dataset_id: UUID, job_id: str):
with session_factory.create_session() as session:
_get_accessible_dataset(dataset_id, current_tenant_id, current_user, session)
if not dify_config.RBAC_ENABLED and not (
current_user.has_edit_permission or current_user.is_dataset_operator
):
raise Forbidden()
snapshots = KnowledgeFSUpgradeSnapshotService(session_factory.get_session_maker())
try:
job = snapshots.retry(tenant_id=current_tenant_id, job_id=job_id)
except KnowledgeFSUpgradeNotFoundError as error:
raise NotFound(str(error)) from error
except KnowledgeFSUpgradeConflictError as error:
raise Conflict(str(error)) from error
if job.old_dataset_id != str(dataset_id):
raise NotFound("Upgrade job was not found")
_enqueue_upgrade_job(snapshots, tenant_id=current_tenant_id, job_id=job.id)
return dump_response(KnowledgeFSUpgradeRetryResponse, {"id": job.id, "status": "queued"}), 202
def _enqueue_upgrade_job(
snapshots: KnowledgeFSUpgradeSnapshotService,
*,
tenant_id: str,
job_id: str,
) -> None:
from tasks.knowledge_fs_upgrade_tasks import run_knowledge_fs_upgrade
task_id = str(uuid4())
if not snapshots.claim_enqueue(tenant_id=tenant_id, job_id=job_id, task_id=task_id):
return
try:
run_knowledge_fs_upgrade.apply_async(kwargs={"job_id": job_id}, task_id=task_id)
except Exception:
snapshots.release_enqueue_claim(tenant_id=tenant_id, job_id=job_id, task_id=task_id)
raise
@console_ns.route("/datasets/<uuid:dataset_id>/use-check")
class DatasetUseCheckApi(Resource):
@console_ns.doc("check_dataset_use")

View File

@ -32,6 +32,7 @@ NON_CORE_COVERAGE_ALLOWLIST = frozenset(
"api/migrations/versions/2026_08_10_1200-7c1e9a4b2d60_add_knowledge_fs_staged_uploads.py",
"api/migrations/versions/2026_08_13_1200-9d4e6f8a1b2c_add_knowledge_fs_space_tag_bindings.py",
"api/migrations/versions/2026_08_17_1200-4f8b2c7d9e10_add_knowledge_fs_icon_background.py",
"api/migrations/versions/2026_08_17_1200-f3a8c1d7e920_add_knowledge_fs_upgrade_jobs.py",
}
)
HUNK_HEADER = re.compile(r"^@@ -\d+(?:,\d+)? \+(\d+)(?:,\d+)? @@")
@ -155,6 +156,8 @@ def is_core_coverage_path(path: str) -> bool:
"api/controllers/openapi/knowledge_fs.py",
"api/events/event_handlers/sync_knowledge_fs_bindings_when_app_published_workflow_updated.py",
"api/extensions/ext_knowledge_fs_observability.py",
"api/services/dataset_knowledge_fs_upgrade_file_lease.py",
"api/services/dataset_knowledge_fs_upgrade_service.py",
"api/services/knowledge_fs_capability.py",
}:
return True

View File

@ -179,6 +179,7 @@ def init_app(app: DifyApp) -> Celery:
"tasks.workflow_run_archive_download_tasks", # workflow-run archive download preparation
"tasks.knowledge_fs_initial_source_preview_tasks", # datasource previews use the standard dataset queue
"tasks.knowledge_fs_failed_retrieval_tasks", # best-effort Workflow quality capture uses dataset workers
"tasks.knowledge_fs_upgrade_tasks", # legacy Dataset upgrades use a dedicated queue and worker
]
day = dify_config.CELERY_BEAT_SCHEDULER_TIME
@ -197,6 +198,10 @@ def init_app(app: DifyApp) -> Celery:
"task": "tasks.knowledge_fs_lifecycle_tasks.cleanup_knowledge_fs_staged_uploads",
"schedule": timedelta(seconds=dify_config.KNOWLEDGE_FS_LIFECYCLE_POLL_INTERVAL_SECONDS),
}
beat_schedule["knowledge_fs_upgrade_file_cleanup"] = {
"task": "tasks.knowledge_fs_upgrade_tasks.cleanup_deferred_knowledge_fs_upgrade_files",
"schedule": timedelta(seconds=dify_config.KNOWLEDGE_FS_LIFECYCLE_POLL_INTERVAL_SECONDS),
}
if dify_config.ENABLE_CONVERSATION_CLEANUP_TASK:
imports.append("tasks.delete_conversation_task")
beat_schedule["conversation_cleanup_sweeper"] = {

View File

@ -1,6 +1,6 @@
{
"schemaVersion": 5,
"subtreeTree": "61759d6a25df9fa4746b6cfd0e795da27cfbb388",
"subtreeTree": "a964e23f1d42d13e959b120192fc74df3d755aa1",
"openapiSha256": "3a712231fa850c4f5151bc283205da9062086a0b693f0d1ab01c2c526323f018",
"capabilityV2AuthManifestSha256": "fc0a47e23cce12544882f0298522b4933002e892b84ce1815df7e81d36a7a0c7",
"capabilityV2AuthTestVectorSha256": "ae0de37b1ff05c40f905cf17a7b410d8971acacf64db07d5ee3d6fecfa559ce3",

View File

@ -0,0 +1,172 @@
"""add KnowledgeFS legacy Dataset upgrade jobs
Revision ID: f3a8c1d7e920
Revises: 4f8b2c7d9e10
Create Date: 2026-08-17 12:00:00.000000
"""
import sqlalchemy as sa
from alembic import op
import models
revision = "f3a8c1d7e920"
down_revision = "4f8b2c7d9e10"
branch_labels = None
depends_on = None
def upgrade() -> None:
op.create_table(
"knowledge_fs_upgrade_jobs",
sa.Column("tenant_id", models.types.StringUUID(), nullable=False),
sa.Column("old_dataset_id", models.types.StringUUID(), nullable=False),
sa.Column("requested_by_account_id", models.types.StringUUID(), nullable=False),
sa.Column("owner_account_id", models.types.StringUUID(), nullable=False),
sa.Column("idempotency_key", sa.String(length=255), nullable=False),
sa.Column("snapshot_at", sa.DateTime(), nullable=False),
sa.Column("config_snapshot", sa.JSON(), nullable=False),
sa.Column("permission_snapshot", sa.JSON(), nullable=False),
sa.Column("app_binding_snapshot", sa.JSON(), nullable=False),
sa.Column("tag_ids_snapshot", sa.JSON(), nullable=False),
sa.Column("status", sa.String(length=16), server_default=sa.text("'queued'"), nullable=False),
sa.Column("stage", sa.String(length=32), server_default=sa.text("'validating'"), nullable=False),
sa.Column("new_control_space_id", models.types.StringUUID(), nullable=True),
sa.Column("resolved_configuration", sa.JSON(), nullable=True),
sa.Column("total_documents", sa.Integer(), server_default=sa.text("0"), nullable=False),
sa.Column("completed_documents", sa.Integer(), server_default=sa.text("0"), nullable=False),
sa.Column("total_sources", sa.Integer(), server_default=sa.text("0"), nullable=False),
sa.Column("completed_sources", sa.Integer(), server_default=sa.text("0"), nullable=False),
sa.Column("attempt_count", sa.Integer(), server_default=sa.text("0"), nullable=False),
sa.Column("celery_task_id", sa.String(length=255), nullable=True),
sa.Column("last_error_code", sa.String(length=128), nullable=True),
sa.Column("last_error_message", models.types.LongText(), nullable=True),
sa.Column("completed_at", sa.DateTime(), nullable=True),
sa.Column("id", models.types.StringUUID(), nullable=False),
sa.Column("created_at", sa.DateTime(), server_default=sa.text("CURRENT_TIMESTAMP"), nullable=False),
sa.Column("updated_at", sa.DateTime(), server_default=sa.text("CURRENT_TIMESTAMP"), nullable=False),
sa.CheckConstraint("attempt_count >= 0", name="kfs_upgrade_job_attempt_count_ck"),
sa.CheckConstraint("completed_documents >= 0", name="kfs_upgrade_job_document_done_ck"),
sa.CheckConstraint("total_documents >= 0", name="kfs_upgrade_job_document_total_ck"),
sa.CheckConstraint("completed_sources >= 0", name="kfs_upgrade_job_source_done_ck"),
sa.CheckConstraint("total_sources >= 0", name="kfs_upgrade_job_source_total_ck"),
sa.ForeignKeyConstraint(
["tenant_id"], ["tenants.id"], name="kfs_upgrade_job_workspace_fk", ondelete="RESTRICT"
),
sa.ForeignKeyConstraint(
["tenant_id", "new_control_space_id"],
["knowledge_fs_control_spaces.tenant_id", "knowledge_fs_control_spaces.id"],
name="kfs_upgrade_job_space_fk",
ondelete="RESTRICT",
),
sa.PrimaryKeyConstraint("id", name="kfs_upgrade_job_pkey"),
sa.UniqueConstraint("tenant_id", "idempotency_key", name="kfs_upgrade_job_idempotency_uq"),
)
op.create_index(
"kfs_upgrade_job_dataset_created_idx",
"knowledge_fs_upgrade_jobs",
["tenant_id", "old_dataset_id", "created_at"],
)
op.create_index("kfs_upgrade_job_status_updated_idx", "knowledge_fs_upgrade_jobs", ["status", "updated_at"])
op.create_table(
"knowledge_fs_upgrade_documents",
sa.Column("job_id", models.types.StringUUID(), nullable=False),
sa.Column("tenant_id", models.types.StringUUID(), nullable=False),
sa.Column("old_document_id", models.types.StringUUID(), nullable=False),
sa.Column("name", sa.String(length=255), nullable=False),
sa.Column("data_source_type", sa.String(length=32), nullable=False),
sa.Column("data_source_info", sa.JSON(), nullable=False),
sa.Column("metadata_snapshot", sa.JSON(), nullable=False),
sa.Column("desired_enabled", sa.Boolean(), nullable=False),
sa.Column("legacy_archived", sa.Boolean(), nullable=False),
sa.Column("legacy_indexing_status", sa.String(length=32), nullable=False),
sa.Column("legacy_display_status", sa.String(length=32), nullable=True),
sa.Column("old_upload_file_id", models.types.StringUUID(), nullable=True),
sa.Column("source_key", sa.String(length=255), nullable=True),
sa.Column("status", sa.String(length=16), server_default=sa.text("'pending'"), nullable=False),
sa.Column("staged_upload_id", models.types.StringUUID(), nullable=True),
sa.Column("new_document_asset_id", models.types.StringUUID(), nullable=True),
sa.Column("new_logical_document_id", models.types.StringUUID(), nullable=True),
sa.Column("compilation_job_id", models.types.StringUUID(), nullable=True),
sa.Column("state_reconcile_attempt_count", sa.Integer(), server_default=sa.text("0"), nullable=False),
sa.Column("state_reconciled_at", sa.DateTime(), nullable=True),
sa.Column("state_reconcile_error", models.types.LongText(), nullable=True),
sa.Column("last_error_code", sa.String(length=128), nullable=True),
sa.Column("last_error_message", models.types.LongText(), nullable=True),
sa.Column("id", models.types.StringUUID(), nullable=False),
sa.Column("created_at", sa.DateTime(), server_default=sa.text("CURRENT_TIMESTAMP"), nullable=False),
sa.Column("updated_at", sa.DateTime(), server_default=sa.text("CURRENT_TIMESTAMP"), nullable=False),
sa.CheckConstraint(
"state_reconcile_attempt_count >= 0",
name="kfs_upgrade_document_reconcile_attempt_ck",
),
sa.ForeignKeyConstraint(
["job_id"], ["knowledge_fs_upgrade_jobs.id"], name="kfs_upgrade_document_job_fk", ondelete="CASCADE"
),
sa.PrimaryKeyConstraint("id", name="kfs_upgrade_document_pkey"),
sa.UniqueConstraint("job_id", "old_document_id", name="kfs_upgrade_document_identity_uq"),
)
op.create_index("kfs_upgrade_document_dispatch_idx", "knowledge_fs_upgrade_documents", ["job_id", "status", "id"])
op.create_table(
"knowledge_fs_upgrade_sources",
sa.Column("job_id", models.types.StringUUID(), nullable=False),
sa.Column("tenant_id", models.types.StringUUID(), nullable=False),
sa.Column("source_key", sa.String(length=255), nullable=False),
sa.Column("payload_snapshot", sa.JSON(), nullable=False),
sa.Column("status", sa.String(length=16), server_default=sa.text("'pending'"), nullable=False),
sa.Column("new_connection_id", models.types.StringUUID(), nullable=True),
sa.Column("new_source_id", models.types.StringUUID(), nullable=True),
sa.Column("initial_sync_task_id", models.types.StringUUID(), nullable=True),
sa.Column("last_error_code", sa.String(length=128), nullable=True),
sa.Column("last_error_message", models.types.LongText(), nullable=True),
sa.Column("id", models.types.StringUUID(), nullable=False),
sa.Column("created_at", sa.DateTime(), server_default=sa.text("CURRENT_TIMESTAMP"), nullable=False),
sa.Column("updated_at", sa.DateTime(), server_default=sa.text("CURRENT_TIMESTAMP"), nullable=False),
sa.ForeignKeyConstraint(
["job_id"], ["knowledge_fs_upgrade_jobs.id"], name="kfs_upgrade_source_job_fk", ondelete="CASCADE"
),
sa.PrimaryKeyConstraint("id", name="kfs_upgrade_source_pkey"),
sa.UniqueConstraint("job_id", "source_key", name="kfs_upgrade_source_identity_uq"),
)
op.create_index("kfs_upgrade_source_dispatch_idx", "knowledge_fs_upgrade_sources", ["job_id", "status", "id"])
op.create_table(
"knowledge_fs_upgrade_file_leases",
sa.Column("job_id", models.types.StringUUID(), nullable=False),
sa.Column("old_upload_file_id", models.types.StringUUID(), nullable=False),
sa.Column("status", sa.String(length=16), server_default=sa.text("'active'"), nullable=False),
sa.Column("expires_at", sa.DateTime(), nullable=False),
sa.Column("released_at", sa.DateTime(), nullable=True),
sa.Column("cleanup_requested_at", sa.DateTime(), nullable=True),
sa.Column("id", models.types.StringUUID(), nullable=False),
sa.Column("created_at", sa.DateTime(), server_default=sa.text("CURRENT_TIMESTAMP"), nullable=False),
sa.Column("updated_at", sa.DateTime(), server_default=sa.text("CURRENT_TIMESTAMP"), nullable=False),
sa.ForeignKeyConstraint(
["job_id"],
["knowledge_fs_upgrade_jobs.id"],
name="kfs_upgrade_file_lease_job_fk",
ondelete="CASCADE",
),
sa.PrimaryKeyConstraint("id", name="kfs_upgrade_file_lease_pkey"),
sa.UniqueConstraint("job_id", "old_upload_file_id", name="kfs_upgrade_file_lease_identity_uq"),
)
op.create_index(
"kfs_upgrade_file_lease_active_idx",
"knowledge_fs_upgrade_file_leases",
["old_upload_file_id", "status", "expires_at"],
)
def downgrade() -> None:
op.drop_index("kfs_upgrade_file_lease_active_idx", table_name="knowledge_fs_upgrade_file_leases")
op.drop_table("knowledge_fs_upgrade_file_leases")
op.drop_index("kfs_upgrade_source_dispatch_idx", table_name="knowledge_fs_upgrade_sources")
op.drop_table("knowledge_fs_upgrade_sources")
op.drop_index("kfs_upgrade_document_dispatch_idx", table_name="knowledge_fs_upgrade_documents")
op.drop_table("knowledge_fs_upgrade_documents")
op.drop_index("kfs_upgrade_job_status_updated_idx", table_name="knowledge_fs_upgrade_jobs")
op.drop_index("kfs_upgrade_job_dataset_created_idx", table_name="knowledge_fs_upgrade_jobs")
op.drop_table("knowledge_fs_upgrade_jobs")

View File

@ -90,6 +90,14 @@ from .knowledge_fs import (
KnowledgeFSSpaceTagBinding,
KnowledgeFSStagedUpload,
KnowledgeFSStagedUploadStatus,
KnowledgeFSUpgradeDocument,
KnowledgeFSUpgradeFileLease,
KnowledgeFSUpgradeFileLeaseStatus,
KnowledgeFSUpgradeItemStatus,
KnowledgeFSUpgradeJob,
KnowledgeFSUpgradeJobStatus,
KnowledgeFSUpgradeSource,
KnowledgeFSUpgradeStage,
)
from .knowledge_fs_cleanup import (
KnowledgeFSCleanupAuthorization,
@ -307,6 +315,14 @@ __all__ = [
"KnowledgeFSSpaceTagBinding",
"KnowledgeFSStagedUpload",
"KnowledgeFSStagedUploadStatus",
"KnowledgeFSUpgradeDocument",
"KnowledgeFSUpgradeFileLease",
"KnowledgeFSUpgradeFileLeaseStatus",
"KnowledgeFSUpgradeItemStatus",
"KnowledgeFSUpgradeJob",
"KnowledgeFSUpgradeJobStatus",
"KnowledgeFSUpgradeSource",
"KnowledgeFSUpgradeStage",
"KnowledgeFSWorkspaceCutoverLedger",
"KnowledgeFSWorkspaceCutoverPhase",
"LoadBalancingModelConfig",

View File

@ -211,6 +211,36 @@ class KnowledgeFSStagedUploadStatus(StrEnum):
EXPIRED = "expired"
class KnowledgeFSUpgradeJobStatus(StrEnum):
QUEUED = "queued"
RUNNING = "running"
SUCCEEDED = "succeeded"
FAILED = "failed"
class KnowledgeFSUpgradeStage(StrEnum):
VALIDATING = "validating"
WAITING_FOR_SPACE = "waiting_for_space"
CREATING_SOURCES = "creating_sources"
SUBMITTING_DOCUMENTS = "submitting_documents"
MIGRATING_ACCESS = "migrating_access"
FINALIZING = "finalizing"
COMPLETED = "completed"
class KnowledgeFSUpgradeItemStatus(StrEnum):
PENDING = "pending"
PROCESSING = "processing"
SUCCEEDED = "succeeded"
FAILED = "failed"
class KnowledgeFSUpgradeFileLeaseStatus(StrEnum):
ACTIVE = "active"
RELEASED = "released"
EXPIRED = "expired"
class KnowledgeFSControlSpace(DefaultFieldsDCMixin, TypeBase):
"""Dify product resource registered to at most one KnowledgeFS Space."""
@ -813,6 +843,183 @@ class KnowledgeFSLifecycleOutbox(DefaultFieldsDCMixin, TypeBase):
retain_until: Mapped[datetime | None] = mapped_column(DateTime, nullable=True, default=None)
class KnowledgeFSUpgradeJob(DefaultFieldsDCMixin, TypeBase):
"""Immutable legacy Dataset snapshot and resumable upgrade progress."""
__tablename__ = "knowledge_fs_upgrade_jobs"
__table_args__ = (
sa.PrimaryKeyConstraint("id", name="kfs_upgrade_job_pkey"),
UniqueConstraint("tenant_id", "idempotency_key", name="kfs_upgrade_job_idempotency_uq"),
sa.ForeignKeyConstraint(
["tenant_id"],
["tenants.id"],
name="kfs_upgrade_job_workspace_fk",
ondelete="RESTRICT",
),
sa.ForeignKeyConstraint(
["tenant_id", "new_control_space_id"],
["knowledge_fs_control_spaces.tenant_id", "knowledge_fs_control_spaces.id"],
name="kfs_upgrade_job_space_fk",
ondelete="RESTRICT",
),
Index("kfs_upgrade_job_dataset_created_idx", "tenant_id", "old_dataset_id", "created_at"),
Index("kfs_upgrade_job_status_updated_idx", "status", "updated_at"),
sa.CheckConstraint("total_documents >= 0", name=sa.schema.conv("kfs_upgrade_job_document_total_ck")),
sa.CheckConstraint("completed_documents >= 0", name=sa.schema.conv("kfs_upgrade_job_document_done_ck")),
sa.CheckConstraint("total_sources >= 0", name=sa.schema.conv("kfs_upgrade_job_source_total_ck")),
sa.CheckConstraint("completed_sources >= 0", name=sa.schema.conv("kfs_upgrade_job_source_done_ck")),
sa.CheckConstraint("attempt_count >= 0", name=sa.schema.conv("kfs_upgrade_job_attempt_count_ck")),
)
tenant_id: Mapped[str] = mapped_column(StringUUID, nullable=False)
old_dataset_id: Mapped[str] = mapped_column(StringUUID, nullable=False)
requested_by_account_id: Mapped[str] = mapped_column(StringUUID, nullable=False)
owner_account_id: Mapped[str] = mapped_column(StringUUID, nullable=False)
idempotency_key: Mapped[str] = mapped_column(String(255), nullable=False)
snapshot_at: Mapped[datetime] = mapped_column(DateTime, nullable=False)
config_snapshot: Mapped[dict[str, object]] = mapped_column(sa.JSON, nullable=False)
permission_snapshot: Mapped[dict[str, object]] = mapped_column(sa.JSON, nullable=False)
app_binding_snapshot: Mapped[list[dict[str, object]]] = mapped_column(sa.JSON, nullable=False)
tag_ids_snapshot: Mapped[list[str]] = mapped_column(sa.JSON, nullable=False)
status: Mapped[KnowledgeFSUpgradeJobStatus] = mapped_column(
EnumText(KnowledgeFSUpgradeJobStatus, length=16),
nullable=False,
server_default=sa.text("'queued'"),
default=KnowledgeFSUpgradeJobStatus.QUEUED,
)
stage: Mapped[KnowledgeFSUpgradeStage] = mapped_column(
EnumText(KnowledgeFSUpgradeStage, length=32),
nullable=False,
server_default=sa.text("'validating'"),
default=KnowledgeFSUpgradeStage.VALIDATING,
)
new_control_space_id: Mapped[str | None] = mapped_column(StringUUID, nullable=True, default=None)
resolved_configuration: Mapped[dict[str, object] | None] = mapped_column(sa.JSON, nullable=True, default=None)
total_documents: Mapped[int] = mapped_column(sa.Integer, nullable=False, default=0, server_default=sa.text("0"))
completed_documents: Mapped[int] = mapped_column(sa.Integer, nullable=False, default=0, server_default=sa.text("0"))
total_sources: Mapped[int] = mapped_column(sa.Integer, nullable=False, default=0, server_default=sa.text("0"))
completed_sources: Mapped[int] = mapped_column(sa.Integer, nullable=False, default=0, server_default=sa.text("0"))
attempt_count: Mapped[int] = mapped_column(sa.Integer, nullable=False, default=0, server_default=sa.text("0"))
celery_task_id: Mapped[str | None] = mapped_column(String(255), nullable=True, default=None)
last_error_code: Mapped[str | None] = mapped_column(String(128), nullable=True, default=None)
last_error_message: Mapped[str | None] = mapped_column(LongText, nullable=True, default=None)
completed_at: Mapped[datetime | None] = mapped_column(DateTime, nullable=True, default=None)
class KnowledgeFSUpgradeDocument(DefaultFieldsDCMixin, TypeBase):
"""One legacy document captured at click time, independent of later Dataset updates."""
__tablename__ = "knowledge_fs_upgrade_documents"
__table_args__ = (
sa.PrimaryKeyConstraint("id", name="kfs_upgrade_document_pkey"),
UniqueConstraint("job_id", "old_document_id", name="kfs_upgrade_document_identity_uq"),
sa.ForeignKeyConstraint(
["job_id"],
["knowledge_fs_upgrade_jobs.id"],
name="kfs_upgrade_document_job_fk",
ondelete="CASCADE",
),
Index("kfs_upgrade_document_dispatch_idx", "job_id", "status", "id"),
sa.CheckConstraint(
"state_reconcile_attempt_count >= 0",
name=sa.schema.conv("kfs_upgrade_document_reconcile_attempt_ck"),
),
)
job_id: Mapped[str] = mapped_column(StringUUID, nullable=False)
tenant_id: Mapped[str] = mapped_column(StringUUID, nullable=False)
old_document_id: Mapped[str] = mapped_column(StringUUID, nullable=False)
name: Mapped[str] = mapped_column(String(255), nullable=False)
data_source_type: Mapped[str] = mapped_column(String(32), nullable=False)
data_source_info: Mapped[dict[str, object]] = mapped_column(sa.JSON, nullable=False)
metadata_snapshot: Mapped[dict[str, object]] = mapped_column(sa.JSON, nullable=False)
desired_enabled: Mapped[bool] = mapped_column(sa.Boolean, nullable=False)
legacy_archived: Mapped[bool] = mapped_column(sa.Boolean, nullable=False)
legacy_indexing_status: Mapped[str] = mapped_column(String(32), nullable=False)
legacy_display_status: Mapped[str | None] = mapped_column(String(32), nullable=True, default=None)
old_upload_file_id: Mapped[str | None] = mapped_column(StringUUID, nullable=True, default=None)
source_key: Mapped[str | None] = mapped_column(String(255), nullable=True, default=None)
status: Mapped[KnowledgeFSUpgradeItemStatus] = mapped_column(
EnumText(KnowledgeFSUpgradeItemStatus, length=16),
nullable=False,
server_default=sa.text("'pending'"),
default=KnowledgeFSUpgradeItemStatus.PENDING,
)
staged_upload_id: Mapped[str | None] = mapped_column(StringUUID, nullable=True, default=None)
new_document_asset_id: Mapped[str | None] = mapped_column(StringUUID, nullable=True, default=None)
new_logical_document_id: Mapped[str | None] = mapped_column(StringUUID, nullable=True, default=None)
compilation_job_id: Mapped[str | None] = mapped_column(StringUUID, nullable=True, default=None)
state_reconcile_attempt_count: Mapped[int] = mapped_column(
sa.Integer, nullable=False, server_default=sa.text("0"), default=0
)
state_reconciled_at: Mapped[datetime | None] = mapped_column(DateTime, nullable=True, default=None)
state_reconcile_error: Mapped[str | None] = mapped_column(LongText, nullable=True, default=None)
last_error_code: Mapped[str | None] = mapped_column(String(128), nullable=True, default=None)
last_error_message: Mapped[str | None] = mapped_column(LongText, nullable=True, default=None)
class KnowledgeFSUpgradeSource(DefaultFieldsDCMixin, TypeBase):
"""One deduplicated Source definition derived from the Dataset snapshot."""
__tablename__ = "knowledge_fs_upgrade_sources"
__table_args__ = (
sa.PrimaryKeyConstraint("id", name="kfs_upgrade_source_pkey"),
UniqueConstraint("job_id", "source_key", name="kfs_upgrade_source_identity_uq"),
sa.ForeignKeyConstraint(
["job_id"],
["knowledge_fs_upgrade_jobs.id"],
name="kfs_upgrade_source_job_fk",
ondelete="CASCADE",
),
Index("kfs_upgrade_source_dispatch_idx", "job_id", "status", "id"),
)
job_id: Mapped[str] = mapped_column(StringUUID, nullable=False)
tenant_id: Mapped[str] = mapped_column(StringUUID, nullable=False)
source_key: Mapped[str] = mapped_column(String(255), nullable=False)
payload_snapshot: Mapped[dict[str, object]] = mapped_column(sa.JSON, nullable=False)
status: Mapped[KnowledgeFSUpgradeItemStatus] = mapped_column(
EnumText(KnowledgeFSUpgradeItemStatus, length=16),
nullable=False,
server_default=sa.text("'pending'"),
default=KnowledgeFSUpgradeItemStatus.PENDING,
)
new_connection_id: Mapped[str | None] = mapped_column(StringUUID, nullable=True, default=None)
new_source_id: Mapped[str | None] = mapped_column(StringUUID, nullable=True, default=None)
initial_sync_task_id: Mapped[str | None] = mapped_column(StringUUID, nullable=True, default=None)
last_error_code: Mapped[str | None] = mapped_column(String(128), nullable=True, default=None)
last_error_message: Mapped[str | None] = mapped_column(LongText, nullable=True, default=None)
class KnowledgeFSUpgradeFileLease(DefaultFieldsDCMixin, TypeBase):
"""Temporary physical-retention lease for a legacy uploaded source file."""
__tablename__ = "knowledge_fs_upgrade_file_leases"
__table_args__ = (
sa.PrimaryKeyConstraint("id", name="kfs_upgrade_file_lease_pkey"),
UniqueConstraint("job_id", "old_upload_file_id", name="kfs_upgrade_file_lease_identity_uq"),
sa.ForeignKeyConstraint(
["job_id"],
["knowledge_fs_upgrade_jobs.id"],
name="kfs_upgrade_file_lease_job_fk",
ondelete="CASCADE",
),
Index("kfs_upgrade_file_lease_active_idx", "old_upload_file_id", "status", "expires_at"),
)
job_id: Mapped[str] = mapped_column(StringUUID, nullable=False)
old_upload_file_id: Mapped[str] = mapped_column(StringUUID, nullable=False)
expires_at: Mapped[datetime] = mapped_column(DateTime, nullable=False)
status: Mapped[KnowledgeFSUpgradeFileLeaseStatus] = mapped_column(
EnumText(KnowledgeFSUpgradeFileLeaseStatus, length=16),
nullable=False,
server_default=sa.text("'active'"),
default=KnowledgeFSUpgradeFileLeaseStatus.ACTIVE,
)
released_at: Mapped[datetime | None] = mapped_column(DateTime, nullable=True, default=None)
cleanup_requested_at: Mapped[datetime | None] = mapped_column(DateTime, nullable=True, default=None)
__all__ = [
"AppKnowledgeFSSpaceJoin",
"KnowledgeFSAllowedActions",
@ -850,4 +1057,12 @@ __all__ = [
"KnowledgeFSScoreThresholdIntentPayload",
"KnowledgeFSStagedUpload",
"KnowledgeFSStagedUploadStatus",
"KnowledgeFSUpgradeDocument",
"KnowledgeFSUpgradeFileLease",
"KnowledgeFSUpgradeFileLeaseStatus",
"KnowledgeFSUpgradeItemStatus",
"KnowledgeFSUpgradeJob",
"KnowledgeFSUpgradeJobStatus",
"KnowledgeFSUpgradeSource",
"KnowledgeFSUpgradeStage",
]

View File

@ -0,0 +1,193 @@
"""Physical-retention checks for legacy files captured by an upgrade snapshot."""
from __future__ import annotations
import json
from collections.abc import Iterable
from datetime import datetime
from json import JSONDecodeError
import sqlalchemy as sa
from sqlalchemy.orm import Session, aliased, sessionmaker
from extensions.ext_storage import storage
from libs.datetime_utils import naive_utc_now
from models.dataset import Document
from models.knowledge_fs import KnowledgeFSUpgradeFileLease, KnowledgeFSUpgradeFileLeaseStatus
from models.model import UploadFile
def active_upgrade_file_ids(
session: Session,
upload_file_ids: Iterable[str],
*,
now: datetime | None = None,
) -> set[str]:
"""Return source files whose physical deletion must wait for migration."""
normalized_ids = frozenset(str(upload_file_id) for upload_file_id in upload_file_ids if upload_file_id)
if not normalized_ids:
return set()
return set(
session.scalars(
sa.select(KnowledgeFSUpgradeFileLease.old_upload_file_id).where(
KnowledgeFSUpgradeFileLease.old_upload_file_id.in_(normalized_ids),
KnowledgeFSUpgradeFileLease.status == KnowledgeFSUpgradeFileLeaseStatus.ACTIVE,
KnowledgeFSUpgradeFileLease.expires_at > (now or naive_utc_now()),
)
)
)
def reserve_upgrade_file_cleanup(
session: Session,
upload_file_ids: Iterable[str],
*,
now: datetime | None = None,
) -> set[str]:
"""Persist cleanup intent for files whose active migration lease blocks deletion."""
requested_at = now or naive_utc_now()
normalized_ids = frozenset(str(upload_file_id) for upload_file_id in upload_file_ids if upload_file_id)
if not normalized_ids:
return set()
leases = list(
session.scalars(
sa.select(KnowledgeFSUpgradeFileLease)
.where(
KnowledgeFSUpgradeFileLease.old_upload_file_id.in_(normalized_ids),
KnowledgeFSUpgradeFileLease.status == KnowledgeFSUpgradeFileLeaseStatus.ACTIVE,
KnowledgeFSUpgradeFileLease.expires_at > requested_at,
)
.with_for_update()
)
)
for lease in leases:
if lease.cleanup_requested_at is None:
lease.cleanup_requested_at = requested_at
return {lease.old_upload_file_id for lease in leases}
def release_upgrade_file_lease(
session: Session,
*,
job_id: str,
upload_file_id: str,
now: datetime | None = None,
) -> bool:
"""Release one lease and fulfill deferred cleanup after the last protection ends."""
released_at = now or naive_utc_now()
leases = list(
session.scalars(
sa.select(KnowledgeFSUpgradeFileLease)
.where(KnowledgeFSUpgradeFileLease.old_upload_file_id == upload_file_id)
.with_for_update()
)
)
current = next((lease for lease in leases if lease.job_id == job_id), None)
if current is None:
return False
current.status = KnowledgeFSUpgradeFileLeaseStatus.RELEASED
current.released_at = released_at
cleanup_requested = any(lease.cleanup_requested_at is not None for lease in leases)
still_protected = any(
lease.status == KnowledgeFSUpgradeFileLeaseStatus.ACTIVE and lease.expires_at > released_at for lease in leases
)
if not cleanup_requested or still_protected or _legacy_document_references_file(session, upload_file_id):
return False
upload_file = session.get(UploadFile, upload_file_id)
if upload_file is None:
return False
storage.delete(upload_file.key)
session.delete(upload_file)
return True
def cleanup_deferred_upgrade_files(
session_maker: sessionmaker[Session],
*,
limit: int = 100,
now: datetime | None = None,
) -> int:
"""Retry bounded cleanup requests after abandoned migration leases expire."""
cleanup_at = now or naive_utc_now()
with session_maker() as session:
candidate = aliased(KnowledgeFSUpgradeFileLease)
blocking = aliased(KnowledgeFSUpgradeFileLease)
upload_file_ids = list(
session.scalars(
sa.select(candidate.old_upload_file_id)
.where(
candidate.cleanup_requested_at.is_not(None),
~sa.exists(
sa.select(blocking.id).where(
blocking.old_upload_file_id == candidate.old_upload_file_id,
blocking.status == KnowledgeFSUpgradeFileLeaseStatus.ACTIVE,
blocking.expires_at > cleanup_at,
)
),
)
.distinct()
.limit(limit)
)
)
cleaned = 0
for upload_file_id in upload_file_ids:
with session_maker.begin() as session:
leases = list(
session.scalars(
sa.select(KnowledgeFSUpgradeFileLease)
.where(KnowledgeFSUpgradeFileLease.old_upload_file_id == upload_file_id)
.with_for_update()
)
)
for lease in leases:
if lease.status == KnowledgeFSUpgradeFileLeaseStatus.ACTIVE and lease.expires_at <= cleanup_at:
lease.status = KnowledgeFSUpgradeFileLeaseStatus.EXPIRED
if any(
lease.status == KnowledgeFSUpgradeFileLeaseStatus.ACTIVE and lease.expires_at > cleanup_at
for lease in leases
):
continue
if _legacy_document_references_file(session, upload_file_id):
for lease in leases:
lease.cleanup_requested_at = None
continue
upload_file = session.get(UploadFile, upload_file_id)
if upload_file is None:
for lease in leases:
lease.cleanup_requested_at = None
continue
storage.delete(upload_file.key)
session.delete(upload_file)
cleaned += 1
return cleaned
def _legacy_document_references_file(session: Session, upload_file_id: str) -> bool:
candidates = session.scalars(
sa.select(Document.data_source_info).where(
Document.data_source_type == "upload_file",
Document.data_source_info.contains(upload_file_id),
)
)
for candidate in candidates:
if not candidate:
continue
try:
payload = json.loads(candidate)
except (JSONDecodeError, TypeError):
continue
if isinstance(payload, dict) and str(payload.get("upload_file_id") or "") == upload_file_id:
return True
return False
__all__ = [
"active_upgrade_file_ids",
"cleanup_deferred_upgrade_files",
"release_upgrade_file_lease",
"reserve_upgrade_file_cleanup",
]

File diff suppressed because it is too large Load Diff

View File

@ -17,6 +17,8 @@ from models.knowledge_fs import (
KnowledgeFSControlSpacePermissionRole,
KnowledgeFSControlSpaceState,
KnowledgeFSControlSpaceVisibility,
KnowledgeFSUpgradeJobStatus,
KnowledgeFSUpgradeStage,
)
from services.knowledge_fs.product_operations import KnowledgeFSProductPermission
@ -294,6 +296,27 @@ class KnowledgeFSSpaceCreatePayload(BaseModel):
return self
class KnowledgeFSUpgradeJobResponse(ResponseModel):
id: str
old_dataset_id: str
new_control_space_id: str | None = None
status: KnowledgeFSUpgradeJobStatus
stage: KnowledgeFSUpgradeStage
snapshot_at: datetime
total_documents: int = Field(ge=0)
completed_documents: int = Field(ge=0)
total_sources: int = Field(ge=0)
completed_sources: int = Field(ge=0)
last_error_code: str | None = None
last_error_message: str | None = None
completed_at: datetime | None = None
class KnowledgeFSUpgradeRetryResponse(ResponseModel):
id: str
status: Literal["queued"] = "queued"
class KnowledgeFSSpaceUpdatePayload(BaseModel):
name: str | None = Field(default=None, min_length=1, max_length=40)
icon: KnowledgeFSIconIdentity | None = None

View File

@ -13,6 +13,7 @@ from core.tools.utils.web_reader_tool import get_image_upload_file_ids
from extensions.ext_storage import storage
from models.dataset import Dataset, DatasetMetadataBinding, DocumentSegment
from models.model import UploadFile
from services.dataset_knowledge_fs_upgrade_file_lease import reserve_upgrade_file_cleanup
from tasks.refresh_billing_vector_space_task import schedule_billing_vector_space_refresh
logger = logging.getLogger(__name__)
@ -47,10 +48,11 @@ def batch_clean_document_task(
segment_ids: list[str] = []
total_image_upload_file_ids: list[str] = []
dataset_tenant_id: str | None = None
deletable_file_ids = list(file_ids)
try:
# ============ Step 1: Query segment and file data (short read-only transaction) ============
with session_factory.create_session() as session:
with session_factory.create_session() as session, session.begin():
# Get segments info
segments = session.scalars(
select(DocumentSegment).where(DocumentSegment.document_id.in_(document_ids))
@ -74,7 +76,15 @@ def batch_clean_document_task(
# Query storage keys for document files
if file_ids:
files = session.scalars(select(UploadFile).where(UploadFile.id.in_(file_ids))).all()
leased_file_ids = reserve_upgrade_file_cleanup(session, file_ids)
deletable_file_ids = [file_id for file_id in file_ids if file_id not in leased_file_ids]
if leased_file_ids:
logger.info(
"Keep %d source files while KnowledgeFS upgrade leases are active, dataset_id=%s",
len(leased_file_ids),
dataset_id,
)
files = session.scalars(select(UploadFile).where(UploadFile.id.in_(deletable_file_ids))).all()
storage_keys_to_delete.extend([f.key for f in files if f and f.key])
# ============ Step 2: Clean vector index (external service, fresh session for dataset) ============
@ -182,17 +192,17 @@ def batch_clean_document_task(
)
# ============ Step 6: Delete document-associated files (separate short transaction) ============
if file_ids:
if deletable_file_ids:
try:
with session_factory.create_session() as session:
stmt = delete(UploadFile).where(UploadFile.id.in_(file_ids))
stmt = delete(UploadFile).where(UploadFile.id.in_(deletable_file_ids))
session.execute(stmt)
session.commit()
except Exception:
logger.exception(
"Failed to delete document UploadFile records for dataset_id: %s, file_ids: %s",
dataset_id,
file_ids,
deletable_file_ids,
)
# ============ Step 7: Delete storage files (I/O operations, no DB transaction) ============

View File

@ -24,6 +24,7 @@ from models.dataset import (
)
from models.model import UploadFile
from models.workflow import Workflow
from services.dataset_knowledge_fs_upgrade_file_lease import reserve_upgrade_file_cleanup
from tasks.refresh_billing_vector_space_task import schedule_billing_vector_space_refresh
logger = logging.getLogger(__name__)
@ -181,11 +182,19 @@ def clean_dataset_task(
if data_source_info and "upload_file_id" in data_source_info:
file_id = data_source_info["upload_file_id"]
file_ids.append(file_id)
files = session.scalars(select(UploadFile).where(UploadFile.id.in_(file_ids))).all()
leased_file_ids = reserve_upgrade_file_cleanup(session, file_ids)
deletable_file_ids = [file_id for file_id in file_ids if file_id not in leased_file_ids]
if leased_file_ids:
logger.info(
"Keep %d source files while KnowledgeFS upgrade leases are active, dataset_id=%s",
len(leased_file_ids),
dataset_id,
)
files = session.scalars(select(UploadFile).where(UploadFile.id.in_(deletable_file_ids))).all()
for file in files:
storage.delete(file.key)
file_delete_stmt = delete(UploadFile).where(UploadFile.id.in_(file_ids))
file_delete_stmt = delete(UploadFile).where(UploadFile.id.in_(deletable_file_ids))
session.execute(file_delete_stmt)
session.commit()

View File

@ -11,6 +11,7 @@ from core.tools.utils.web_reader_tool import get_image_upload_file_ids
from extensions.ext_storage import storage
from models.dataset import Dataset, DatasetMetadataBinding, DocumentSegment, SegmentAttachmentBinding
from models.model import UploadFile
from services.dataset_knowledge_fs_upgrade_file_lease import reserve_upgrade_file_cleanup
from tasks.refresh_billing_vector_space_task import schedule_billing_vector_space_refresh
logger = logging.getLogger(__name__)
@ -125,6 +126,14 @@ def clean_document_task(
)
with session_factory.create_session() as session, session.begin():
if file_id:
if file_id in reserve_upgrade_file_cleanup(session, [file_id]):
logger.info(
"Keep source file while KnowledgeFS upgrade lease is active, file_id=%s, document_id=%s",
file_id,
document_id,
)
file_id = None
if file_id:
file = session.scalar(select(UploadFile).where(UploadFile.id == file_id).limit(1))
if file:

View File

@ -56,6 +56,16 @@ class _DatasourceBinding:
provider_kind: str
@dataclass(frozen=True)
class KnowledgeFSInitialSourceSubmission:
"""Source creation result; the first import is deliberately best-effort."""
connection_id: str
source_id: str
workflow_id: str | None
workflow_error: str | None = None
class KnowledgeFSInitialSourceNotReadyError(RuntimeError):
"""The Space, connection, or Source workflow is progressing and should be retried."""
@ -512,6 +522,140 @@ def start_initial_source_import(
return workflow.id
def submit_initial_source_for_upgrade(
*,
tenant_id: str,
account_id: str,
control_space_id: str,
operation_id: str,
payload: KnowledgeFSInitialSourcePayload,
) -> KnowledgeFSInitialSourceSubmission:
"""Create and commit one Source without waiting for its import workflow.
Upgrade success is based on the independently usable Source existing. The
selected import is submitted when possible, but a failure is returned as a
warning because users can retry it from the new KnowledgeFS task surface.
"""
session_maker = session_factory.get_session_maker()
with session_maker() as session:
control_space = SQLAlchemyKnowledgeFSControlSpaceRepository(session).get(
tenant_id=tenant_id,
control_space_id=control_space_id,
)
if control_space is None:
raise RuntimeError("KnowledgeFS control-space was not found")
if control_space.state is KnowledgeFSControlSpaceState.PROVISIONING:
raise KnowledgeFSInitialSourceNotReadyError("KnowledgeFS Space is still provisioning")
if control_space.state is not KnowledgeFSControlSpaceState.ACTIVE or control_space.knowledge_space_id is None:
raise RuntimeError(f"KnowledgeFS Space cannot accept a Source in state {control_space.state.value}")
facade = get_knowledge_fs_runtime(session_maker).facade
request_id = _request_id(operation_id=operation_id, payload=payload)
source = _find_initial_source(
facade=facade,
tenant_id=tenant_id,
account_id=account_id,
control_space_id=control_space_id,
request_id=request_id,
)
if source is None:
binding = _binding(payload)
credential_id, credential_name = _find_credential(
session_maker=session_maker,
tenant_id=tenant_id,
account_id=account_id,
binding=binding,
)
connection = _find_or_create_connection(
facade=facade,
tenant_id=tenant_id,
account_id=account_id,
control_space_id=control_space_id,
binding=binding,
credential_id=credential_id,
credential_name=credential_name,
)
source = facade.create_source(
tenant_id=tenant_id,
account_id=account_id,
control_space_id=control_space_id,
payload=_source_payload(
payload=payload,
binding=binding,
connection_id=connection.id,
request_id=request_id,
),
)
elif source.connection_id is None:
raise RuntimeError("Initial Source has no connection")
if source.status != "active" or source.metadata.get("preview") is not False:
source = facade.update_source(
tenant_id=tenant_id,
account_id=account_id,
control_space_id=control_space_id,
source_id=source.id,
payload=KnowledgeFSSourceUpdatePayload(
expectedVersion=source.version,
metadata={**source.metadata, "preview": False, "upgradeJobId": operation_id},
status="active",
),
)
try:
current_policy = facade.get_source_sync_policy(
tenant_id=tenant_id,
account_id=account_id,
control_space_id=control_space_id,
source_id=source.id,
)
expected_revision = current_policy.revision
except KnowledgeFSProductResourceNotFoundError:
expected_revision = 0
facade.update_source_sync_policy(
tenant_id=tenant_id,
account_id=account_id,
control_space_id=control_space_id,
source_id=source.id,
payload=_sync_policy_payload(
payload=payload,
expected_revision=expected_revision,
source_version=source.version,
),
)
workflow_id: str | None = None
workflow_error: str | None = None
try:
workflow = _start_workflow(
facade=facade,
tenant_id=tenant_id,
account_id=account_id,
control_space_id=control_space_id,
source_id=source.id,
request_id=request_id,
payload=payload,
)
workflow_id = workflow.id
except Exception as exc:
workflow_error = type(exc).__name__
logger.warning(
"KnowledgeFS upgrade Source was created but its first import was not submitted",
extra={
"control_space_id": control_space_id,
"error_code": workflow_error,
"source_id": source.id,
},
)
return KnowledgeFSInitialSourceSubmission(
connection_id=source.connection_id or "",
source_id=source.id,
workflow_id=workflow_id,
workflow_error=workflow_error,
)
def start_initial_website_source_import(
*,
tenant_id: str,
@ -637,8 +781,10 @@ def import_initial_website_source(
__all__ = [
"KnowledgeFSInitialSourceSubmission",
"import_initial_source",
"import_initial_website_source",
"start_initial_source_import",
"start_initial_website_source_import",
"submit_initial_source_for_upgrade",
]

View File

@ -0,0 +1,93 @@
"""Dedicated-queue execution for legacy Dataset upgrades."""
from __future__ import annotations
import logging
from celery import shared_task
from core.db.session_factory import session_factory
from services.dataset_knowledge_fs_upgrade_file_lease import cleanup_deferred_upgrade_files
from services.dataset_knowledge_fs_upgrade_service import (
KnowledgeFSUpgradeDocumentReconciler,
KnowledgeFSUpgradeNotReadyError,
KnowledgeFSUpgradeRunner,
)
KNOWLEDGE_FS_UPGRADE_QUEUE = "knowledge_fs_upgrade"
logger = logging.getLogger(__name__)
@shared_task(
bind=True,
queue=KNOWLEDGE_FS_UPGRADE_QUEUE,
max_retries=360,
default_retry_delay=5,
)
def run_knowledge_fs_upgrade(self, *, job_id: str) -> None:
"""Run one checkpoint and enqueue the next only after it is committed."""
runner = KnowledgeFSUpgradeRunner(session_factory.get_session_maker())
try:
has_more = runner.run_next(job_id=job_id, celery_task_id=self.request.id)
except KnowledgeFSUpgradeNotReadyError as error:
if self.request.retries >= self.max_retries:
runner.fail(job_id=job_id, error=error)
logger.exception(
"KnowledgeFS Dataset upgrade exhausted provisioning retries",
extra={"upgrade_job_id": job_id},
)
raise
raise self.retry(exc=error)
except Exception as error:
runner.fail(job_id=job_id, error=error)
logger.exception("KnowledgeFS Dataset upgrade failed", extra={"upgrade_job_id": job_id})
raise
if has_more:
run_knowledge_fs_upgrade.apply_async(kwargs={"job_id": job_id})
else:
reconcile_knowledge_fs_upgrade_documents.apply_async(kwargs={"job_id": job_id})
@shared_task(
bind=True,
queue=KNOWLEDGE_FS_UPGRADE_QUEUE,
max_retries=10_080,
default_retry_delay=60,
)
def reconcile_knowledge_fs_upgrade_documents(self, *, job_id: str) -> None:
"""Eventually apply click-time metadata and availability without gating migration success."""
reconciler = KnowledgeFSUpgradeDocumentReconciler(session_factory.get_session_maker())
try:
remaining = reconciler.reconcile(job_id=job_id)
except Exception as error:
logger.warning(
"KnowledgeFS Dataset upgrade document reconciliation is not ready",
extra={"upgrade_job_id": job_id},
exc_info=True,
)
raise self.retry(exc=error)
if remaining:
raise self.retry(exc=KnowledgeFSUpgradeNotReadyError(f"{remaining} migrated documents are not visible yet"))
@shared_task(queue=KNOWLEDGE_FS_UPGRADE_QUEUE)
def cleanup_deferred_knowledge_fs_upgrade_files() -> int:
"""Delete orphaned source files after abandoned upgrade leases expire."""
return cleanup_deferred_upgrade_files(session_factory.get_session_maker())
def enqueue_knowledge_fs_upgrade(*, job_id: str) -> str:
result = run_knowledge_fs_upgrade.apply_async(kwargs={"job_id": job_id})
return str(result.id)
__all__ = [
"KNOWLEDGE_FS_UPGRADE_QUEUE",
"cleanup_deferred_knowledge_fs_upgrade_files",
"enqueue_knowledge_fs_upgrade",
"reconcile_knowledge_fs_upgrade_documents",
"run_knowledge_fs_upgrade",
]

View File

@ -23,6 +23,8 @@ from controllers.console.datasets.datasets import (
DatasetErrorDocs,
DatasetIndexingEstimateApi,
DatasetIndexingStatusApi,
DatasetKnowledgeFSUpgradeApi,
DatasetKnowledgeFSUpgradeJobApi,
DatasetListApi,
DatasetPermissionUserListApi,
DatasetQueryApi,
@ -45,6 +47,7 @@ from extensions.storage.storage_type import StorageType
from models.account import Account, TenantAccountRole
from models.dataset import Dataset, DatasetQuery, Document
from models.enums import CreatorUserRole, DataSourceType, DocumentCreatedFrom, IndexingStatus
from models.knowledge_fs import KnowledgeFSUpgradeJobStatus, KnowledgeFSUpgradeStage
from models.model import ApiToken, App, AppMode, IconType, UploadFile
from services.dataset_ref_service import DatasetRef
from services.dataset_service import DatasetPermissionService, DatasetService
@ -827,6 +830,144 @@ class TestDatasetApiDelete:
method(api, MagicMock(), user, dataset_id)
class TestDatasetKnowledgeFSUpgradeApi:
@staticmethod
def _job(dataset_id: str, **overrides):
values = {
"id": "upgrade-job-1",
"old_dataset_id": dataset_id,
"new_control_space_id": None,
"status": KnowledgeFSUpgradeJobStatus.QUEUED,
"stage": KnowledgeFSUpgradeStage.VALIDATING,
"snapshot_at": datetime.datetime(2026, 8, 17, tzinfo=datetime.UTC),
"total_documents": 1,
"completed_documents": 0,
"total_sources": 0,
"completed_sources": 0,
"last_error_code": None,
"last_error_message": None,
"completed_at": None,
}
values.update(overrides)
return SimpleNamespace(**values)
def test_create_snapshots_and_enqueues_without_remote_work_in_request(self, app: Flask):
dataset_id = "123e4567-e89b-12d3-a456-426614174000"
dataset = make_dataset(id=dataset_id, tenant_id="tenant-1")
user = make_account()
snapshots = MagicMock()
snapshots.create.return_value = self._job(dataset_id)
session_context = MagicMock()
session_context.__enter__.return_value = MagicMock()
api = DatasetKnowledgeFSUpgradeApi()
method = unwrap(api.post)
with (
app.test_request_context(
f"/datasets/{dataset_id}/knowledge-fs-upgrades",
headers={"Idempotency-Key": "upgrade-request-1"},
),
patch("controllers.console.datasets.datasets.dify_config.KNOWLEDGE_FS_ENABLED", True),
patch("controllers.console.datasets.datasets.dify_config.RBAC_ENABLED", False),
patch("controllers.console.datasets.datasets.session_factory.create_session", return_value=session_context),
patch("controllers.console.datasets.datasets.session_factory.get_session_maker", return_value="maker"),
patch.object(DatasetService, "get_dataset_for_tenant", return_value=dataset),
patch.object(DatasetService, "check_dataset_permission"),
patch(
"controllers.console.datasets.datasets.KnowledgeFSUpgradeSnapshotService",
return_value=snapshots,
),
patch("controllers.console.datasets.datasets._enqueue_upgrade_job") as enqueue,
):
response, status = method(api, "tenant-1", user, dataset_id)
assert status == 202
assert response["id"] == "upgrade-job-1"
snapshots.create.assert_called_once_with(
tenant_id="tenant-1",
dataset_id=dataset_id,
requested_by_account_id=user.id,
idempotency_key="upgrade-request-1",
)
enqueue.assert_called_once_with(snapshots, tenant_id="tenant-1", job_id="upgrade-job-1")
def test_status_uses_legacy_dataset_permission_in_community_edition(self, app: Flask):
dataset_id = "123e4567-e89b-12d3-a456-426614174000"
dataset = make_dataset(id=dataset_id, tenant_id="tenant-1")
api = DatasetKnowledgeFSUpgradeJobApi()
method = unwrap(api.get)
session_context = MagicMock()
session_context.__enter__.return_value = MagicMock()
with (
app.test_request_context(f"/datasets/{dataset_id}/knowledge-fs-upgrades/job-1"),
patch("controllers.console.datasets.datasets.dify_config.RBAC_ENABLED", False),
patch("controllers.console.datasets.datasets.session_factory.create_session", return_value=session_context),
patch.object(DatasetService, "get_dataset_for_tenant", return_value=dataset),
patch.object(
DatasetService,
"check_dataset_permission",
side_effect=services.errors.account.NoPermissionError("no access"),
),
patch("controllers.console.datasets.datasets.KnowledgeFSUpgradeSnapshotService") as snapshots,
pytest.raises(Forbidden, match="no access"),
):
method(api, "tenant-1", make_account(), dataset_id, "job-1")
snapshots.assert_not_called()
def test_status_relies_on_rbac_decorator_when_enterprise_rbac_is_enabled(self, app: Flask):
dataset_id = "123e4567-e89b-12d3-a456-426614174000"
dataset = make_dataset(id=dataset_id, tenant_id="tenant-1")
snapshots = MagicMock()
snapshots.get.return_value = self._job(dataset_id)
session_context = MagicMock()
session_context.__enter__.return_value = MagicMock()
api = DatasetKnowledgeFSUpgradeJobApi()
method = unwrap(api.get)
with (
app.test_request_context(f"/datasets/{dataset_id}/knowledge-fs-upgrades/job-1"),
patch("controllers.console.datasets.datasets.dify_config.RBAC_ENABLED", True),
patch("controllers.console.datasets.datasets.session_factory.create_session", return_value=session_context),
patch("controllers.console.datasets.datasets.session_factory.get_session_maker", return_value="maker"),
patch.object(DatasetService, "get_dataset_for_tenant", return_value=dataset),
patch.object(DatasetService, "check_dataset_permission") as legacy_permission,
patch(
"controllers.console.datasets.datasets.KnowledgeFSUpgradeSnapshotService",
return_value=snapshots,
),
):
response = method(api, "tenant-1", make_account(), dataset_id, "job-1")
assert response["id"] == "upgrade-job-1"
legacy_permission.assert_not_called()
def test_retry_rejects_a_job_from_another_dataset(self, app: Flask):
dataset_id = "123e4567-e89b-12d3-a456-426614174000"
dataset = make_dataset(id=dataset_id, tenant_id="tenant-1")
snapshots = MagicMock()
snapshots.retry.return_value = self._job("223e4567-e89b-12d3-a456-426614174000")
session_context = MagicMock()
session_context.__enter__.return_value = MagicMock()
api = DatasetKnowledgeFSUpgradeJobApi()
method = unwrap(api.post)
with (
app.test_request_context(f"/datasets/{dataset_id}/knowledge-fs-upgrades/job-1"),
patch("controllers.console.datasets.datasets.dify_config.RBAC_ENABLED", False),
patch("controllers.console.datasets.datasets.session_factory.create_session", return_value=session_context),
patch("controllers.console.datasets.datasets.session_factory.get_session_maker", return_value="maker"),
patch.object(DatasetService, "get_dataset_for_tenant", return_value=dataset),
patch.object(DatasetService, "check_dataset_permission"),
patch(
"controllers.console.datasets.datasets.KnowledgeFSUpgradeSnapshotService",
return_value=snapshots,
),
patch("controllers.console.datasets.datasets._enqueue_upgrade_job") as enqueue,
pytest.raises(NotFound, match="Upgrade job was not found"),
):
method(api, "tenant-1", make_account(), dataset_id, "job-1")
enqueue.assert_not_called()
class TestDatasetUseCheckApi:
@pytest.mark.parametrize("is_using", [True, False])
def test_get_use_check(self, app: Flask, is_using: bool):

View File

@ -1085,7 +1085,7 @@ def test_service_profile_rejects_cross_control_space_before_facade_io() -> None:
def test_product_modules_do_not_import_dify_dataset_or_document_services() -> None:
paths = [
*_API_ROOT.glob("services/knowledge_fs/*.py"),
*(path for path in _API_ROOT.glob("services/knowledge_fs/*.py") if not path.name.startswith("upgrade_")),
*_API_ROOT.glob("controllers/console/knowledge_fs/*.py"),
*_API_ROOT.glob("controllers/service_api/knowledge_fs/*.py"),
]

View File

@ -33,6 +33,8 @@ from dev.check_knowledge_fs_coverage import (
"api/extensions/ext_knowledge_fs_observability.py",
"api/models/knowledge_fs.py",
"api/repositories/sqlalchemy_knowledge_fs_cutover_repository.py",
"api/services/dataset_knowledge_fs_upgrade_file_lease.py",
"api/services/dataset_knowledge_fs_upgrade_service.py",
"api/services/knowledge_fs/runtime.py",
"api/services/knowledge_fs_capability.py",
"api/tasks/knowledge_fs_lifecycle_tasks.py",
@ -121,6 +123,7 @@ def test_core_coverage_rejects_unclassified_knowledge_fs_production_files(tmp_pa
"api/migrations/versions/2026_07_21_1600-e5a7c9b2d416_add_knowledge_fs_cleanup_completion.py",
"api/migrations/versions/2026_08_13_1200-9d4e6f8a1b2c_add_knowledge_fs_space_tag_bindings.py",
"api/migrations/versions/2026_08_17_1200-4f8b2c7d9e10_add_knowledge_fs_icon_background.py",
"api/migrations/versions/2026_08_17_1200-f3a8c1d7e920_add_knowledge_fs_upgrade_jobs.py",
],
)
def test_core_coverage_allows_explicit_non_core_migrations(tmp_path: Path, migration: str) -> None:

View File

@ -66,6 +66,10 @@ def test_celery_registers_initial_source_task_when_knowledge_fs_lifecycle_is_rea
"task": "tasks.knowledge_fs_lifecycle_tasks.cleanup_knowledge_fs_staged_uploads",
"schedule": timedelta(seconds=2),
}
assert celery_app.conf["beat_schedule"]["knowledge_fs_upgrade_file_cleanup"] == {
"task": "tasks.knowledge_fs_upgrade_tasks.cleanup_deferred_knowledge_fs_upgrade_files",
"schedule": timedelta(seconds=2),
}
with (
patch("extensions.ext_celery.dify_config", config),

View File

@ -0,0 +1,62 @@
from __future__ import annotations
import importlib.util
from pathlib import Path
import sqlalchemy as sa
from alembic.migration import MigrationContext
from alembic.operations import Operations
from models.knowledge_fs import (
KnowledgeFSUpgradeDocument,
KnowledgeFSUpgradeFileLease,
KnowledgeFSUpgradeJob,
KnowledgeFSUpgradeSource,
)
_MIGRATION_PATH = (
Path(__file__).resolve().parents[3]
/ "migrations/versions/2026_08_17_1200-f3a8c1d7e920_add_knowledge_fs_upgrade_jobs.py"
)
_MODELS = (
KnowledgeFSUpgradeJob,
KnowledgeFSUpgradeDocument,
KnowledgeFSUpgradeSource,
KnowledgeFSUpgradeFileLease,
)
def _load_migration_module():
spec = importlib.util.spec_from_file_location("knowledge_fs_upgrade_jobs", _MIGRATION_PATH)
if spec is None or spec.loader is None:
raise RuntimeError("failed to load migration module")
module = importlib.util.module_from_spec(spec)
spec.loader.exec_module(module)
return module
def _run_step(module: object, engine: sa.Engine, step_name: str) -> None:
with engine.begin() as connection:
operations = Operations(MigrationContext.configure(connection))
original_op = module.op
module.op = operations
try:
getattr(module, step_name)()
finally:
module.op = original_op
def test_upgrade_schema_matches_upgrade_models_and_downgrades_cleanly() -> None:
engine = sa.create_engine("sqlite:///:memory:")
module = _load_migration_module()
_run_step(module, engine, "upgrade")
inspector = sa.inspect(engine)
assert set(inspector.get_table_names()) == {model.__tablename__ for model in _MODELS}
for model in _MODELS:
migrated_columns = {column["name"] for column in inspector.get_columns(model.__tablename__)}
assert migrated_columns == set(model.__table__.columns.keys())
_run_step(module, engine, "downgrade")
assert sa.inspect(engine).get_table_names() == []

View File

@ -0,0 +1,295 @@
from __future__ import annotations
from datetime import timedelta
from unittest.mock import patch
from sqlalchemy.orm import Session, sessionmaker
from extensions.storage.storage_type import StorageType
from libs.datetime_utils import naive_utc_now
from models.dataset import Document
from models.enums import CreatorUserRole, DataSourceType, DocumentCreatedFrom
from models.knowledge_fs import (
KnowledgeFSUpgradeFileLease,
KnowledgeFSUpgradeFileLeaseStatus,
KnowledgeFSUpgradeJob,
)
from models.model import UploadFile
from services.dataset_knowledge_fs_upgrade_file_lease import (
active_upgrade_file_ids,
cleanup_deferred_upgrade_files,
release_upgrade_file_lease,
reserve_upgrade_file_cleanup,
)
_TENANT_ID = "00000000-0000-0000-0000-000000000001"
_ACCOUNT_ID = "00000000-0000-0000-0000-000000000002"
_DATASET_ID = "00000000-0000-0000-0000-000000000003"
_ACTIVE_FILE_ID = "00000000-0000-0000-0000-000000000004"
_EXPIRED_FILE_ID = "00000000-0000-0000-0000-000000000005"
def test_only_active_unexpired_source_files_are_protected(sqlite_session_factory: sessionmaker[Session]) -> None:
now = naive_utc_now()
job = KnowledgeFSUpgradeJob(
tenant_id=_TENANT_ID,
old_dataset_id=_DATASET_ID,
requested_by_account_id=_ACCOUNT_ID,
owner_account_id=_ACCOUNT_ID,
idempotency_key="upgrade-lease-test",
snapshot_at=now,
config_snapshot={},
permission_snapshot={},
app_binding_snapshot=[],
tag_ids_snapshot=[],
)
with sqlite_session_factory.begin() as session:
session.add(job)
session.flush()
session.add_all(
[
KnowledgeFSUpgradeFileLease(
job_id=job.id,
old_upload_file_id=_ACTIVE_FILE_ID,
expires_at=now + timedelta(minutes=1),
),
KnowledgeFSUpgradeFileLease(
job_id=job.id,
old_upload_file_id=_EXPIRED_FILE_ID,
expires_at=now - timedelta(seconds=1),
status=KnowledgeFSUpgradeFileLeaseStatus.ACTIVE,
),
]
)
with sqlite_session_factory() as session:
assert active_upgrade_file_ids(
session,
[_ACTIVE_FILE_ID, _EXPIRED_FILE_ID],
now=now,
) == {_ACTIVE_FILE_ID}
def test_cleanup_request_is_persisted_on_every_active_lease(
sqlite_session_factory: sessionmaker[Session],
) -> None:
now = naive_utc_now()
jobs = [
KnowledgeFSUpgradeJob(
tenant_id=_TENANT_ID,
old_dataset_id=_DATASET_ID,
requested_by_account_id=_ACCOUNT_ID,
owner_account_id=_ACCOUNT_ID,
idempotency_key=f"upgrade-cleanup-request-{index}",
snapshot_at=now,
config_snapshot={},
permission_snapshot={},
app_binding_snapshot=[],
tag_ids_snapshot=[],
)
for index in range(2)
]
with sqlite_session_factory.begin() as session:
session.add_all(jobs)
session.flush()
leases = [
KnowledgeFSUpgradeFileLease(
job_id=job.id,
old_upload_file_id=_ACTIVE_FILE_ID,
expires_at=now + timedelta(minutes=1),
)
for job in jobs
]
session.add_all(leases)
with sqlite_session_factory.begin() as session:
assert reserve_upgrade_file_cleanup(session, [_ACTIVE_FILE_ID], now=now) == {_ACTIVE_FILE_ID}
with sqlite_session_factory() as session:
persisted = list(session.query(KnowledgeFSUpgradeFileLease).filter_by(old_upload_file_id=_ACTIVE_FILE_ID))
assert len(persisted) == 2
assert all(lease.cleanup_requested_at == now for lease in persisted)
def test_last_lease_release_deletes_deferred_orphan_file(
sqlite_session_factory: sessionmaker[Session],
) -> None:
now = naive_utc_now()
job = KnowledgeFSUpgradeJob(
tenant_id=_TENANT_ID,
old_dataset_id=_DATASET_ID,
requested_by_account_id=_ACCOUNT_ID,
owner_account_id=_ACCOUNT_ID,
idempotency_key="upgrade-deferred-file-cleanup",
snapshot_at=now,
config_snapshot={},
permission_snapshot={},
app_binding_snapshot=[],
tag_ids_snapshot=[],
)
upload_file = UploadFile(
tenant_id=_TENANT_ID,
storage_type=StorageType.LOCAL,
key=f"upload_files/{_TENANT_ID}/orphan.txt",
name="orphan.txt",
size=10,
extension="txt",
mime_type="text/plain",
created_by_role=CreatorUserRole.ACCOUNT,
created_by=_ACCOUNT_ID,
created_at=now,
used=False,
)
with sqlite_session_factory.begin() as session:
session.add_all([job, upload_file])
session.flush()
lease = KnowledgeFSUpgradeFileLease(
job_id=job.id,
old_upload_file_id=upload_file.id,
expires_at=now + timedelta(minutes=1),
cleanup_requested_at=now,
)
session.add(lease)
with patch("services.dataset_knowledge_fs_upgrade_file_lease.storage") as storage:
assert cleanup_deferred_upgrade_files(sqlite_session_factory, now=now) == 0
storage.delete.assert_not_called()
with sqlite_session_factory.begin() as session:
assert (
release_upgrade_file_lease(
session,
job_id=job.id,
upload_file_id=upload_file.id,
now=now,
)
is True
)
storage.delete.assert_called_once_with(upload_file.key)
with sqlite_session_factory() as session:
assert session.get(UploadFile, upload_file.id) is None
persisted_lease = session.get(KnowledgeFSUpgradeFileLease, lease.id)
assert persisted_lease is not None
assert persisted_lease.status is KnowledgeFSUpgradeFileLeaseStatus.RELEASED
def test_lease_release_keeps_file_referenced_by_another_legacy_document(
sqlite_session_factory: sessionmaker[Session],
) -> None:
now = naive_utc_now()
job = KnowledgeFSUpgradeJob(
tenant_id=_TENANT_ID,
old_dataset_id=_DATASET_ID,
requested_by_account_id=_ACCOUNT_ID,
owner_account_id=_ACCOUNT_ID,
idempotency_key="upgrade-shared-legacy-file",
snapshot_at=now,
config_snapshot={},
permission_snapshot={},
app_binding_snapshot=[],
tag_ids_snapshot=[],
)
upload_file = UploadFile(
tenant_id=_TENANT_ID,
storage_type=StorageType.LOCAL,
key=f"upload_files/{_TENANT_ID}/shared.txt",
name="shared.txt",
size=10,
extension="txt",
mime_type="text/plain",
created_by_role=CreatorUserRole.ACCOUNT,
created_by=_ACCOUNT_ID,
created_at=now,
used=False,
)
with sqlite_session_factory.begin() as session:
session.add_all([job, upload_file])
session.flush()
document = Document(
tenant_id=_TENANT_ID,
dataset_id="00000000-0000-0000-0000-000000000020",
position=1,
data_source_type=DataSourceType.UPLOAD_FILE,
data_source_info=f'{{"upload_file_id":"{upload_file.id}"}}',
batch="shared-file-test",
name="shared.txt",
created_from=DocumentCreatedFrom.WEB,
created_by=_ACCOUNT_ID,
enabled=True,
archived=False,
indexing_status="completed",
)
lease = KnowledgeFSUpgradeFileLease(
job_id=job.id,
old_upload_file_id=upload_file.id,
expires_at=now + timedelta(minutes=1),
cleanup_requested_at=now,
)
session.add_all([document, lease])
with patch("services.dataset_knowledge_fs_upgrade_file_lease.storage") as storage:
with sqlite_session_factory.begin() as session:
assert (
release_upgrade_file_lease(
session,
job_id=job.id,
upload_file_id=upload_file.id,
now=now,
)
is False
)
storage.delete.assert_not_called()
with sqlite_session_factory() as session:
assert session.get(UploadFile, upload_file.id) is not None
def test_expired_lease_sweeper_deletes_deferred_orphan_file(
sqlite_session_factory: sessionmaker[Session],
) -> None:
now = naive_utc_now()
job = KnowledgeFSUpgradeJob(
tenant_id=_TENANT_ID,
old_dataset_id=_DATASET_ID,
requested_by_account_id=_ACCOUNT_ID,
owner_account_id=_ACCOUNT_ID,
idempotency_key="upgrade-expired-file-cleanup",
snapshot_at=now,
config_snapshot={},
permission_snapshot={},
app_binding_snapshot=[],
tag_ids_snapshot=[],
)
upload_file = UploadFile(
tenant_id=_TENANT_ID,
storage_type=StorageType.LOCAL,
key=f"upload_files/{_TENANT_ID}/expired-orphan.txt",
name="expired-orphan.txt",
size=10,
extension="txt",
mime_type="text/plain",
created_by_role=CreatorUserRole.ACCOUNT,
created_by=_ACCOUNT_ID,
created_at=now,
used=False,
)
with sqlite_session_factory.begin() as session:
session.add_all([job, upload_file])
session.flush()
lease = KnowledgeFSUpgradeFileLease(
job_id=job.id,
old_upload_file_id=upload_file.id,
expires_at=now - timedelta(seconds=1),
cleanup_requested_at=now - timedelta(minutes=1),
)
session.add(lease)
with patch("services.dataset_knowledge_fs_upgrade_file_lease.storage") as storage:
assert cleanup_deferred_upgrade_files(sqlite_session_factory, now=now) == 1
storage.delete.assert_called_once_with(upload_file.key)
with sqlite_session_factory() as session:
assert session.get(UploadFile, upload_file.id) is None
persisted_lease = session.get(KnowledgeFSUpgradeFileLease, lease.id)
assert persisted_lease is not None
assert persisted_lease.status is KnowledgeFSUpgradeFileLeaseStatus.EXPIRED

File diff suppressed because it is too large Load Diff

View File

@ -1,12 +1,16 @@
import uuid
from datetime import UTC, datetime, timedelta
from unittest.mock import patch
import pytest
from sqlalchemy.orm import Session
import tasks.batch_clean_document_task as task_module
from extensions.storage.storage_type import StorageType
from models.dataset import Dataset, DocumentSegment
from models.enums import DataSourceType
from models.enums import CreatorUserRole, DataSourceType
from models.knowledge_fs import KnowledgeFSUpgradeFileLease, KnowledgeFSUpgradeJob
from models.model import UploadFile
from tasks.batch_clean_document_task import batch_clean_document_task
@ -83,3 +87,80 @@ def test_failed_vector_cleanup_does_not_schedule_billing_refresh(cleanup_rows: t
)
schedule_refresh.assert_not_called()
def test_batch_cleanup_keeps_only_the_leased_legacy_source_file(
cleanup_rows: tuple[str, str, str], sqlite_session: Session
) -> None:
dataset_id, document_id, tenant_id = cleanup_rows
account_id = str(uuid.uuid4())
now = datetime.now(UTC).replace(tzinfo=None)
leased_file = UploadFile(
tenant_id=tenant_id,
storage_type=StorageType.LOCAL,
key=f"upload_files/{tenant_id}/leased.txt",
name="leased.txt",
size=10,
extension="txt",
mime_type="text/plain",
created_by_role=CreatorUserRole.ACCOUNT,
created_by=account_id,
created_at=now,
used=False,
)
deletable_file = UploadFile(
tenant_id=tenant_id,
storage_type=StorageType.LOCAL,
key=f"upload_files/{tenant_id}/deletable.txt",
name="deletable.txt",
size=10,
extension="txt",
mime_type="text/plain",
created_by_role=CreatorUserRole.ACCOUNT,
created_by=account_id,
created_at=now,
used=False,
)
job = KnowledgeFSUpgradeJob(
tenant_id=tenant_id,
old_dataset_id=dataset_id,
requested_by_account_id=account_id,
owner_account_id=account_id,
idempotency_key="batch-cleanup-lease-test",
snapshot_at=now,
config_snapshot={},
permission_snapshot={},
app_binding_snapshot=[],
tag_ids_snapshot=[],
)
sqlite_session.add_all([leased_file, deletable_file, job])
sqlite_session.flush()
sqlite_session.add(
KnowledgeFSUpgradeFileLease(
job_id=job.id,
old_upload_file_id=leased_file.id,
expires_at=now + timedelta(hours=1),
)
)
sqlite_session.commit()
leased_file_id = leased_file.id
deletable_file_id = deletable_file.id
deletable_file_key = deletable_file.key
with (
patch("tasks.batch_clean_document_task.get_image_upload_file_ids", return_value=[]),
patch("tasks.batch_clean_document_task.IndexProcessorFactory"),
patch("tasks.batch_clean_document_task.schedule_billing_vector_space_refresh"),
patch("tasks.batch_clean_document_task.storage") as storage,
):
batch_clean_document_task(
document_ids=[document_id],
dataset_id=dataset_id,
doc_form="paragraph",
file_ids=[leased_file_id, deletable_file_id],
)
sqlite_session.expire_all()
assert sqlite_session.get(UploadFile, leased_file_id) is not None
assert sqlite_session.get(UploadFile, deletable_file_id) is None
storage.delete.assert_called_once_with(deletable_file_key)

View File

@ -14,7 +14,7 @@ This module tests the dataset cleanup task functionality including:
import json
import uuid
from collections.abc import Iterator
from datetime import UTC, datetime
from datetime import UTC, datetime, timedelta
from unittest.mock import MagicMock, patch
import pytest
@ -36,6 +36,7 @@ from models.dataset import (
SegmentAttachmentBinding,
)
from models.enums import CreatorUserRole, DataSourceType, DocumentCreatedFrom, IndexingStatus
from models.knowledge_fs import KnowledgeFSUpgradeFileLease, KnowledgeFSUpgradeJob
from models.model import UploadFile
from models.workflow import Workflow, WorkflowType
from tasks.clean_dataset_task import clean_dataset_task
@ -579,3 +580,77 @@ class TestIndexProcessorParameters:
)
schedule_refresh.assert_not_called()
def test_dataset_cleanup_keeps_a_leased_legacy_source_file(
dataset_id: str,
tenant_id: str,
collection_binding_id: str,
orm_session_maker: sessionmaker[Session],
mock_storage: MagicMock,
mock_index_processor_factory: dict[str, MagicMock],
mock_get_image_upload_file_ids: MagicMock,
) -> None:
del mock_index_processor_factory, mock_get_image_upload_file_ids
account_id = str(uuid.uuid4())
now = datetime.now(UTC).replace(tzinfo=None)
upload_file = UploadFile(
tenant_id=tenant_id,
storage_type=StorageType.LOCAL,
key=f"upload_files/{tenant_id}/dataset-source.txt",
name="dataset-source.txt",
size=10,
extension="txt",
mime_type="text/plain",
created_by_role=CreatorUserRole.ACCOUNT,
created_by=account_id,
created_at=now,
used=False,
)
document = Document(
id=str(uuid.uuid4()),
tenant_id=tenant_id,
dataset_id=dataset_id,
position=1,
data_source_type=DataSourceType.UPLOAD_FILE,
data_source_info=json.dumps({"upload_file_id": upload_file.id}),
batch="batch",
name="dataset-source.txt",
created_from=DocumentCreatedFrom.WEB,
created_by=account_id,
indexing_status=IndexingStatus.COMPLETED,
doc_form=IndexStructureType.PARAGRAPH_INDEX,
)
job = KnowledgeFSUpgradeJob(
tenant_id=tenant_id,
old_dataset_id=dataset_id,
requested_by_account_id=account_id,
owner_account_id=account_id,
idempotency_key="dataset-cleanup-lease-test",
snapshot_at=now,
config_snapshot={},
permission_snapshot={},
app_binding_snapshot=[],
tag_ids_snapshot=[],
)
with orm_session_maker.begin() as session:
session.add_all([upload_file, document, job])
session.flush()
session.add(
KnowledgeFSUpgradeFileLease(
job_id=job.id,
old_upload_file_id=upload_file.id,
expires_at=now + timedelta(hours=1),
)
)
_run_clean_dataset(
dataset_id=dataset_id,
tenant_id=tenant_id,
collection_binding_id=collection_binding_id,
)
with orm_session_maker() as session:
assert session.get(Document, document.id) is None
assert session.get(UploadFile, upload_file.id) is not None
mock_storage.delete.assert_not_called()

View File

@ -6,6 +6,7 @@ starts from the production incident shape: the caller has already deleted the
"""
import uuid
from datetime import UTC, datetime, timedelta
from unittest.mock import MagicMock, patch
import pytest
@ -13,6 +14,7 @@ from sqlalchemy import select
from sqlalchemy.orm import Session
import tasks.clean_document_task as clean_document_task_module
from extensions.storage.storage_type import StorageType
from models.dataset import (
Dataset,
DatasetMetadataBinding,
@ -20,7 +22,8 @@ from models.dataset import (
DocumentSegment,
SegmentAttachmentBinding,
)
from models.enums import DataSourceType, DocumentCreatedFrom
from models.enums import CreatorUserRole, DataSourceType, DocumentCreatedFrom
from models.knowledge_fs import KnowledgeFSUpgradeFileLease, KnowledgeFSUpgradeJob
from models.model import UploadFile
from tasks.clean_document_task import clean_document_task
@ -244,6 +247,77 @@ class TestVectorCleanupResilience:
)
schedule_refresh.assert_not_called()
def test_active_upgrade_lease_keeps_the_legacy_source_file(
document_id: str,
dataset_id: str,
tenant_id: str,
sqlite_session: Session,
bind_task_sessions: None,
mock_storage,
mock_index_processor_factory,
) -> None:
del bind_task_sessions, mock_index_processor_factory
_persist_deleted_document_state(
sqlite_session,
document_id=document_id,
dataset_id=dataset_id,
tenant_id=tenant_id,
target_segment_ids=[],
)
account_id = str(uuid.uuid4())
upload_file = UploadFile(
tenant_id=tenant_id,
storage_type=StorageType.LOCAL,
key=f"upload_files/{tenant_id}/source.txt",
name="source.txt",
size=12,
extension="txt",
mime_type="text/plain",
created_by_role=CreatorUserRole.ACCOUNT,
created_by=account_id,
created_at=datetime(2026, 8, 17, tzinfo=UTC),
used=False,
)
now = datetime.now(UTC).replace(tzinfo=None)
job = KnowledgeFSUpgradeJob(
tenant_id=tenant_id,
old_dataset_id=dataset_id,
requested_by_account_id=account_id,
owner_account_id=account_id,
idempotency_key="cleanup-lease-test",
snapshot_at=now,
config_snapshot={},
permission_snapshot={},
app_binding_snapshot=[],
tag_ids_snapshot=[],
)
sqlite_session.add_all([upload_file, job])
sqlite_session.flush()
lease = KnowledgeFSUpgradeFileLease(
job_id=job.id,
old_upload_file_id=upload_file.id,
expires_at=now + timedelta(hours=1),
)
sqlite_session.add(lease)
sqlite_session.commit()
clean_document_task(
document_id=document_id,
dataset_id=dataset_id,
doc_form="paragraph",
file_id=upload_file.id,
)
sqlite_session.expire_all()
assert sqlite_session.get(UploadFile, upload_file.id) is not None
persisted_lease = sqlite_session.get(KnowledgeFSUpgradeFileLease, lease.id)
assert persisted_lease is not None
assert persisted_lease.cleanup_requested_at is not None
mock_storage.delete.assert_not_called()
class TestVectorCleanupSuccessPaths:
def test_vector_cleanup_success_path_remains_unaffected(
self,
document_id: str,

View File

@ -17,6 +17,7 @@ from tasks.knowledge_fs_initial_source_tasks import (
import_initial_website_source,
start_initial_source_import,
start_initial_website_source_import,
submit_initial_source_for_upgrade,
)
_DEFAULT_CREDENTIAL = object()
@ -794,3 +795,37 @@ def test_initial_source_task_does_not_retry_authoritative_missing_resource() ->
)
retry.assert_not_called()
def test_upgrade_source_submission_succeeds_when_first_import_submission_fails() -> None:
facade = _facade()
facade.create_source.return_value = SimpleNamespace(
connection_id="connection-1",
id="source-1",
metadata={"clientRequestId": "initial-website-source:operation-1", "preview": True},
status="disabled",
version=1,
)
facade.update_source.return_value = SimpleNamespace(
connection_id="connection-1",
id="source-1",
metadata={"clientRequestId": "initial-website-source:operation-1", "preview": False},
status="active",
version=2,
)
facade.import_selected_source_crawl.side_effect = RuntimeError("crawl queue unavailable")
with _runtime(facade):
result = submit_initial_source_for_upgrade(
tenant_id="tenant-1",
account_id="account-1",
control_space_id="control-1",
operation_id="operation-1",
payload=_payload(sync_policy="manual"),
)
assert result.connection_id == "connection-1"
assert result.source_id == "source-1"
assert result.workflow_id is None
assert result.workflow_error == "RuntimeError"
facade.update_source_sync_policy.assert_called_once()

View File

@ -0,0 +1,82 @@
from __future__ import annotations
from unittest.mock import MagicMock, patch
import pytest
from services.dataset_knowledge_fs_upgrade_service import KnowledgeFSUpgradeNotReadyError
from tasks.knowledge_fs_upgrade_tasks import (
KNOWLEDGE_FS_UPGRADE_QUEUE,
cleanup_deferred_knowledge_fs_upgrade_files,
reconcile_knowledge_fs_upgrade_documents,
run_knowledge_fs_upgrade,
)
def test_upgrade_tasks_are_pinned_to_the_dedicated_queue() -> None:
assert run_knowledge_fs_upgrade._get_exec_options()["queue"] == KNOWLEDGE_FS_UPGRADE_QUEUE
assert reconcile_knowledge_fs_upgrade_documents._get_exec_options()["queue"] == KNOWLEDGE_FS_UPGRADE_QUEUE
assert cleanup_deferred_knowledge_fs_upgrade_files._get_exec_options()["queue"] == KNOWLEDGE_FS_UPGRADE_QUEUE
def test_deferred_file_cleanup_uses_the_upgrade_session_factory() -> None:
with (
patch("tasks.knowledge_fs_upgrade_tasks.session_factory.get_session_maker", return_value="maker"),
patch("tasks.knowledge_fs_upgrade_tasks.cleanup_deferred_upgrade_files", return_value=3) as cleanup,
):
assert cleanup_deferred_knowledge_fs_upgrade_files.run() == 3
cleanup.assert_called_once_with("maker")
@pytest.mark.parametrize("has_more", [True, False])
def test_worker_enqueues_only_the_expected_follow_up(has_more: bool) -> None:
runner = MagicMock()
runner.run_next.return_value = has_more
with (
patch("tasks.knowledge_fs_upgrade_tasks.session_factory.get_session_maker", return_value="maker"),
patch("tasks.knowledge_fs_upgrade_tasks.KnowledgeFSUpgradeRunner", return_value=runner),
patch.object(run_knowledge_fs_upgrade, "apply_async") as continue_upgrade,
patch.object(reconcile_knowledge_fs_upgrade_documents, "apply_async") as reconcile,
):
run_knowledge_fs_upgrade.run(job_id="job-1")
runner.run_next.assert_called_once()
if has_more:
continue_upgrade.assert_called_once_with(kwargs={"job_id": "job-1"})
reconcile.assert_not_called()
else:
continue_upgrade.assert_not_called()
reconcile.assert_called_once_with(kwargs={"job_id": "job-1"})
def test_worker_marks_parent_failed_when_a_checkpoint_raises() -> None:
runner = MagicMock()
error = RuntimeError("ordinary document upload failed")
runner.run_next.side_effect = error
with (
patch("tasks.knowledge_fs_upgrade_tasks.session_factory.get_session_maker", return_value="maker"),
patch("tasks.knowledge_fs_upgrade_tasks.KnowledgeFSUpgradeRunner", return_value=runner),
pytest.raises(RuntimeError, match="ordinary document upload failed"),
):
run_knowledge_fs_upgrade.run(job_id="job-1")
runner.fail.assert_called_once_with(job_id="job-1", error=error)
def test_worker_marks_parent_failed_when_provisioning_retries_are_exhausted() -> None:
runner = MagicMock()
error = KnowledgeFSUpgradeNotReadyError("Space is still provisioning")
runner.run_next.side_effect = error
run_knowledge_fs_upgrade.push_request(retries=run_knowledge_fs_upgrade.max_retries)
try:
with (
patch("tasks.knowledge_fs_upgrade_tasks.session_factory.get_session_maker", return_value="maker"),
patch("tasks.knowledge_fs_upgrade_tasks.KnowledgeFSUpgradeRunner", return_value=runner),
pytest.raises(KnowledgeFSUpgradeNotReadyError, match="still provisioning"),
):
run_knowledge_fs_upgrade.run(job_id="job-1")
finally:
run_knowledge_fs_upgrade.pop_request()
runner.fail.assert_called_once_with(job_id="job-1", error=error)

View File

@ -346,6 +346,36 @@ services:
- ssrf_proxy_network
- default
# Dedicated worker for click-time Dataset snapshots and KnowledgeFS upgrade orchestration.
# Keep this queue out of the generic worker so large file copies cannot starve existing jobs.
worker_knowledge_fs_upgrade:
<<: *shared-worker-config
image: langgenius/dify-api:1.16.1
environment:
MODE: worker
CELERY_WORKER_QUEUES: knowledge_fs_upgrade
CELERY_WORKER_CONCURRENCY: ${KNOWLEDGE_FS_UPGRADE_WORKER_AMOUNT:-2}
CELERY_PREFETCH_MULTIPLIER: 1
SENTRY_DSN: ${API_SENTRY_DSN:-}
SENTRY_TRACES_SAMPLE_RATE: ${API_SENTRY_TRACES_SAMPLE_RATE:-1.0}
SENTRY_PROFILES_SAMPLE_RATE: ${API_SENTRY_PROFILES_SAMPLE_RATE:-1.0}
depends_on:
init_permissions:
condition: service_completed_successfully
db_postgres:
condition: service_healthy
required: false
db_mysql:
condition: service_healthy
required: false
redis:
condition: service_started
volumes:
- ./volumes/app/storage:/app/api/storage
networks:
- ssrf_proxy_network
- default
# worker_beat service
# Celery beat for scheduling periodic tasks.
worker_beat:

View File

@ -352,6 +352,36 @@ services:
- ssrf_proxy_network
- default
# Dedicated worker for click-time Dataset snapshots and KnowledgeFS upgrade orchestration.
# Keep this queue out of the generic worker so large file copies cannot starve existing jobs.
worker_knowledge_fs_upgrade:
<<: *shared-worker-config
image: langgenius/dify-api:1.16.1
environment:
MODE: worker
CELERY_WORKER_QUEUES: knowledge_fs_upgrade
CELERY_WORKER_CONCURRENCY: ${KNOWLEDGE_FS_UPGRADE_WORKER_AMOUNT:-2}
CELERY_PREFETCH_MULTIPLIER: 1
SENTRY_DSN: ${API_SENTRY_DSN:-}
SENTRY_TRACES_SAMPLE_RATE: ${API_SENTRY_TRACES_SAMPLE_RATE:-1.0}
SENTRY_PROFILES_SAMPLE_RATE: ${API_SENTRY_PROFILES_SAMPLE_RATE:-1.0}
depends_on:
init_permissions:
condition: service_completed_successfully
db_postgres:
condition: service_healthy
required: false
db_mysql:
condition: service_healthy
required: false
redis:
condition: service_started
volumes:
- ./volumes/app/storage:/app/api/storage
networks:
- ssrf_proxy_network
- default
# worker_beat service
# Celery beat for scheduling periodic tasks.
worker_beat:

View File

@ -59,14 +59,14 @@
},
{
"fingerprint": "ac8011331002c601c2ec9c4dd7f2328e4d3bdc511021e21d4f9e3b67f0f96bc6",
"line": 558,
"line": 588,
"path": "api/models/knowledge_fs.py",
"reason": "Fixed database schema identifier; this literal is not credential material.",
"rule": "dify-knowledge-fs-credential"
},
{
"fingerprint": "fe0a852bad3684b159154a873bcfe0e721e771aa1a4bc0d6b5fee9cb47cf17fd",
"line": 596,
"line": 626,
"path": "api/models/knowledge_fs.py",
"reason": "Fixed database schema identifier; this literal is not credential material.",
"rule": "dify-knowledge-fs-credential"

View File

@ -66,6 +66,8 @@ import {
zGetDatasetsByDatasetIdErrorDocsResponse,
zGetDatasetsByDatasetIdIndexingStatusPath,
zGetDatasetsByDatasetIdIndexingStatusResponse,
zGetDatasetsByDatasetIdKnowledgeFsUpgradesByJobIdPath,
zGetDatasetsByDatasetIdKnowledgeFsUpgradesByJobIdResponse,
zGetDatasetsByDatasetIdMetadataPath,
zGetDatasetsByDatasetIdMetadataResponse,
zGetDatasetsByDatasetIdNotionSyncPath,
@ -162,6 +164,10 @@ import {
zPostDatasetsByDatasetIdHitTestingBody,
zPostDatasetsByDatasetIdHitTestingPath,
zPostDatasetsByDatasetIdHitTestingResponse,
zPostDatasetsByDatasetIdKnowledgeFsUpgradesByJobIdPath,
zPostDatasetsByDatasetIdKnowledgeFsUpgradesByJobIdResponse,
zPostDatasetsByDatasetIdKnowledgeFsUpgradesPath,
zPostDatasetsByDatasetIdKnowledgeFsUpgradesResponse,
zPostDatasetsByDatasetIdMetadataBody,
zPostDatasetsByDatasetIdMetadataBuiltInByActionPath,
zPostDatasetsByDatasetIdMetadataBuiltInByActionResponse,
@ -1424,7 +1430,52 @@ export const indexingStatus3 = {
get: get27,
}
export const get28 = oc
.route({
inputStructure: 'detailed',
method: 'GET',
operationId: 'getDatasetsByDatasetIdKnowledgeFsUpgradesByJobId',
path: '/datasets/{dataset_id}/knowledge-fs-upgrades/{job_id}',
tags: ['console'],
})
.input(z.object({ params: zGetDatasetsByDatasetIdKnowledgeFsUpgradesByJobIdPath }))
.output(zGetDatasetsByDatasetIdKnowledgeFsUpgradesByJobIdResponse)
export const post19 = oc
.route({
inputStructure: 'detailed',
method: 'POST',
operationId: 'postDatasetsByDatasetIdKnowledgeFsUpgradesByJobId',
path: '/datasets/{dataset_id}/knowledge-fs-upgrades/{job_id}',
successStatus: 202,
tags: ['console'],
})
.input(z.object({ params: zPostDatasetsByDatasetIdKnowledgeFsUpgradesByJobIdPath }))
.output(zPostDatasetsByDatasetIdKnowledgeFsUpgradesByJobIdResponse)
export const byJobId2 = {
get: get28,
post: post19,
}
export const post20 = oc
.route({
inputStructure: 'detailed',
method: 'POST',
operationId: 'postDatasetsByDatasetIdKnowledgeFsUpgrades',
path: '/datasets/{dataset_id}/knowledge-fs-upgrades',
successStatus: 202,
tags: ['console'],
})
.input(z.object({ params: zPostDatasetsByDatasetIdKnowledgeFsUpgradesPath }))
.output(zPostDatasetsByDatasetIdKnowledgeFsUpgradesResponse)
export const knowledgeFsUpgrades = {
post: post20,
byJobId: byJobId2,
}
export const post21 = oc
.route({
inputStructure: 'detailed',
method: 'POST',
@ -1437,7 +1488,7 @@ export const post19 = oc
.output(zPostDatasetsByDatasetIdMetadataBuiltInByActionResponse)
export const byAction4 = {
post: post19,
post: post21,
}
export const builtIn2 = {
@ -1477,7 +1528,7 @@ export const byMetadataId = {
patch: patch10,
}
export const get28 = oc
export const get29 = oc
.route({
inputStructure: 'detailed',
method: 'GET',
@ -1488,7 +1539,7 @@ export const get28 = oc
.input(z.object({ params: zGetDatasetsByDatasetIdMetadataPath }))
.output(zGetDatasetsByDatasetIdMetadataResponse)
export const post20 = oc
export const post22 = oc
.route({
inputStructure: 'detailed',
method: 'POST',
@ -1506,13 +1557,13 @@ export const post20 = oc
.output(zPostDatasetsByDatasetIdMetadataResponse)
export const metadata4 = {
get: get28,
post: post20,
get: get29,
post: post22,
builtIn: builtIn2,
byMetadataId,
}
export const get29 = oc
export const get30 = oc
.route({
inputStructure: 'detailed',
method: 'GET',
@ -1524,7 +1575,7 @@ export const get29 = oc
.output(zGetDatasetsByDatasetIdNotionSyncResponse)
export const sync2 = {
get: get29,
get: get30,
}
export const notion2 = {
@ -1534,7 +1585,7 @@ export const notion2 = {
/**
* Get dataset permission user list
*/
export const get30 = oc
export const get31 = oc
.route({
description: 'Get dataset permission user list',
inputStructure: 'detailed',
@ -1547,13 +1598,13 @@ export const get30 = oc
.output(zGetDatasetsByDatasetIdPermissionPartUsersResponse)
export const permissionPartUsers = {
get: get30,
get: get31,
}
/**
* Get dataset query history
*/
export const get31 = oc
export const get32 = oc
.route({
description: 'Get dataset query history',
inputStructure: 'detailed',
@ -1566,13 +1617,13 @@ export const get31 = oc
.output(zGetDatasetsByDatasetIdQueriesResponse)
export const queries = {
get: get31,
get: get32,
}
/**
* Get applications related to dataset
*/
export const get32 = oc
export const get33 = oc
.route({
description: 'Get applications related to dataset',
inputStructure: 'detailed',
@ -1585,13 +1636,13 @@ export const get32 = oc
.output(zGetDatasetsByDatasetIdRelatedAppsResponse)
export const relatedApps = {
get: get32,
get: get33,
}
/**
* retry document
*/
export const post21 = oc
export const post23 = oc
.route({
inputStructure: 'detailed',
method: 'POST',
@ -1610,13 +1661,13 @@ export const post21 = oc
.output(zPostDatasetsByDatasetIdRetryResponse)
export const retry = {
post: post21,
post: post23,
}
/**
* Check if dataset is in use
*/
export const get33 = oc
export const get34 = oc
.route({
description: 'Check if dataset is in use',
inputStructure: 'detailed',
@ -1629,7 +1680,7 @@ export const get33 = oc
.output(zGetDatasetsByDatasetIdUseCheckResponse)
export const useCheck2 = {
get: get33,
get: get34,
}
export const delete9 = oc
@ -1647,7 +1698,7 @@ export const delete9 = oc
/**
* Get dataset details
*/
export const get34 = oc
export const get35 = oc
.route({
description: 'Get dataset details',
inputStructure: 'detailed',
@ -1676,7 +1727,7 @@ export const patch11 = oc
export const byDatasetId = {
delete: delete9,
get: get34,
get: get35,
patch: patch11,
apiKeys: apiKeys2,
autoDisableLogs,
@ -1686,6 +1737,7 @@ export const byDatasetId = {
externalHitTesting,
hitTesting,
indexingStatus: indexingStatus3,
knowledgeFsUpgrades,
metadata: metadata4,
notion: notion2,
permissionPartUsers,
@ -1723,7 +1775,7 @@ export const byApiKeyId2 = {
*
* Get all API keys for a dataset
*/
export const get35 = oc
export const get36 = oc
.route({
description: 'Get all API keys for a dataset',
inputStructure: 'detailed',
@ -1741,7 +1793,7 @@ export const get35 = oc
*
* Create a new API key for a dataset
*/
export const post22 = oc
export const post24 = oc
.route({
description: 'Create a new API key for a dataset',
inputStructure: 'detailed',
@ -1756,8 +1808,8 @@ export const post22 = oc
.output(zPostDatasetsByResourceIdApiKeysResponse)
export const apiKeys3 = {
get: get35,
post: post22,
get: get36,
post: post24,
byApiKeyId: byApiKeyId2,
}
@ -1768,7 +1820,7 @@ export const byResourceId = {
/**
* Get list of datasets
*/
export const get36 = oc
export const get37 = oc
.route({
description: 'Get list of datasets',
inputStructure: 'detailed',
@ -1783,7 +1835,7 @@ export const get36 = oc
/**
* Create a new dataset
*/
export const post23 = oc
export const post25 = oc
.route({
description: 'Create a new dataset',
inputStructure: 'detailed',
@ -1797,8 +1849,8 @@ export const post23 = oc
.output(zPostDatasetsResponse)
export const datasets = {
get: get36,
post: post23,
get: get37,
post: post25,
apiBaseInfo,
apiKeys,
batchImportStatus,

View File

@ -513,6 +513,27 @@ export type HitTestingResponse = {
records: Array<HitTestingRecord>
}
export type KnowledgeFsUpgradeJobResponse = {
completed_at?: string | null
completed_documents: number
completed_sources: number
id: string
last_error_code?: string | null
last_error_message?: string | null
new_control_space_id?: string | null
old_dataset_id: string
snapshot_at: string
stage: KnowledgeFsUpgradeStage
status: KnowledgeFsUpgradeJobStatus
total_documents: number
total_sources: number
}
export type KnowledgeFsUpgradeRetryResponse = {
id: string
status?: 'queued'
}
export type DatasetMetadataListResponse = {
built_in_field_enabled: boolean
doc_metadata: Array<DatasetMetadataListItemResponse>
@ -848,6 +869,17 @@ export type HitTestingRecord = {
tsne_position: unknown | null
}
export type KnowledgeFsUpgradeStage =
| 'completed'
| 'creating_sources'
| 'finalizing'
| 'migrating_access'
| 'submitting_documents'
| 'validating'
| 'waiting_for_space'
export type KnowledgeFsUpgradeJobStatus = 'failed' | 'queued' | 'running' | 'succeeded'
export type DatasetMetadataListItemResponse = {
count?: number
id: string
@ -2302,6 +2334,56 @@ export type GetDatasetsByDatasetIdIndexingStatusResponses = {
export type GetDatasetsByDatasetIdIndexingStatusResponse =
GetDatasetsByDatasetIdIndexingStatusResponses[keyof GetDatasetsByDatasetIdIndexingStatusResponses]
export type PostDatasetsByDatasetIdKnowledgeFsUpgradesData = {
body?: never
path: {
dataset_id: string
}
query?: never
url: '/datasets/{dataset_id}/knowledge-fs-upgrades'
}
export type PostDatasetsByDatasetIdKnowledgeFsUpgradesResponses = {
202: KnowledgeFsUpgradeJobResponse
}
export type PostDatasetsByDatasetIdKnowledgeFsUpgradesResponse =
PostDatasetsByDatasetIdKnowledgeFsUpgradesResponses[keyof PostDatasetsByDatasetIdKnowledgeFsUpgradesResponses]
export type GetDatasetsByDatasetIdKnowledgeFsUpgradesByJobIdData = {
body?: never
path: {
dataset_id: string
job_id: string
}
query?: never
url: '/datasets/{dataset_id}/knowledge-fs-upgrades/{job_id}'
}
export type GetDatasetsByDatasetIdKnowledgeFsUpgradesByJobIdResponses = {
200: KnowledgeFsUpgradeJobResponse
}
export type GetDatasetsByDatasetIdKnowledgeFsUpgradesByJobIdResponse =
GetDatasetsByDatasetIdKnowledgeFsUpgradesByJobIdResponses[keyof GetDatasetsByDatasetIdKnowledgeFsUpgradesByJobIdResponses]
export type PostDatasetsByDatasetIdKnowledgeFsUpgradesByJobIdData = {
body?: never
path: {
dataset_id: string
job_id: string
}
query?: never
url: '/datasets/{dataset_id}/knowledge-fs-upgrades/{job_id}'
}
export type PostDatasetsByDatasetIdKnowledgeFsUpgradesByJobIdResponses = {
202: KnowledgeFsUpgradeRetryResponse
}
export type PostDatasetsByDatasetIdKnowledgeFsUpgradesByJobIdResponse =
PostDatasetsByDatasetIdKnowledgeFsUpgradesByJobIdResponses[keyof PostDatasetsByDatasetIdKnowledgeFsUpgradesByJobIdResponses]
export type GetDatasetsByDatasetIdMetadataData = {
body?: never
path: {

View File

@ -244,6 +244,14 @@ export const zExternalHitTestingPayload = z.object({
query: z.string(),
})
/**
* KnowledgeFSUpgradeRetryResponse
*/
export const zKnowledgeFsUpgradeRetryResponse = z.object({
id: z.string(),
status: z.literal('queued').optional().default('queued'),
})
/**
* MetadataArgs
*/
@ -754,6 +762,43 @@ export const zHitTestingQuery = z.object({
content: z.string(),
})
/**
* KnowledgeFSUpgradeStage
*/
export const zKnowledgeFsUpgradeStage = z.enum([
'completed',
'creating_sources',
'finalizing',
'migrating_access',
'submitting_documents',
'validating',
'waiting_for_space',
])
/**
* KnowledgeFSUpgradeJobStatus
*/
export const zKnowledgeFsUpgradeJobStatus = z.enum(['failed', 'queued', 'running', 'succeeded'])
/**
* KnowledgeFSUpgradeJobResponse
*/
export const zKnowledgeFsUpgradeJobResponse = z.object({
completed_at: z.iso.datetime().nullish(),
completed_documents: z.int().gte(0),
completed_sources: z.int().gte(0),
id: z.string(),
last_error_code: z.string().nullish(),
last_error_message: z.string().nullish(),
new_control_space_id: z.string().nullish(),
old_dataset_id: z.string(),
snapshot_at: z.iso.datetime(),
stage: zKnowledgeFsUpgradeStage,
status: zKnowledgeFsUpgradeJobStatus,
total_documents: z.int().gte(0),
total_sources: z.int().gte(0),
})
/**
* DatasetMetadataListItemResponse
*/
@ -2176,6 +2221,37 @@ export const zGetDatasetsByDatasetIdIndexingStatusPath = z.object({
*/
export const zGetDatasetsByDatasetIdIndexingStatusResponse = zDocumentStatusListResponse
export const zPostDatasetsByDatasetIdKnowledgeFsUpgradesPath = z.object({
dataset_id: z.uuid(),
})
/**
* KnowledgeFS Dataset upgrade accepted
*/
export const zPostDatasetsByDatasetIdKnowledgeFsUpgradesResponse = zKnowledgeFsUpgradeJobResponse
export const zGetDatasetsByDatasetIdKnowledgeFsUpgradesByJobIdPath = z.object({
dataset_id: z.uuid(),
job_id: z.string(),
})
/**
* KnowledgeFS Dataset upgrade status
*/
export const zGetDatasetsByDatasetIdKnowledgeFsUpgradesByJobIdResponse =
zKnowledgeFsUpgradeJobResponse
export const zPostDatasetsByDatasetIdKnowledgeFsUpgradesByJobIdPath = z.object({
dataset_id: z.uuid(),
job_id: z.string(),
})
/**
* KnowledgeFS Dataset upgrade retry accepted
*/
export const zPostDatasetsByDatasetIdKnowledgeFsUpgradesByJobIdResponse =
zKnowledgeFsUpgradeRetryResponse
export const zGetDatasetsByDatasetIdMetadataPath = z.object({
dataset_id: z.uuid(),
})

View File

@ -0,0 +1,155 @@
'use client'
import type { KnowledgeFsUpgradeJobResponse } from '@dify/contracts/api/console/datasets/types.gen'
import { Button } from '@langgenius/dify-ui/button'
import { useMutation, useQuery } from '@tanstack/react-query'
import { useState } from 'react'
import { useTranslation } from 'react-i18next'
import { consoleQuery } from '@/service/client'
type Props = {
datasetId: string
disabled: boolean
}
const POLL_INTERVAL = 2_000
const getErrorMessage = (error: unknown) => (error instanceof Error ? error.message : String(error))
const UpgradeProgress = ({
datasetId,
initialJob,
}: {
datasetId: string
initialJob: KnowledgeFsUpgradeJobResponse
}) => {
const { t } = useTranslation()
const [retryError, setRetryError] = useState<string>()
const jobContract = consoleQuery.datasets.byDatasetId.knowledgeFsUpgrades.byJobId
const jobInput = {
params: {
dataset_id: datasetId,
job_id: initialJob.id,
},
}
const {
data: job,
error: statusError,
refetch,
} = useQuery({
...jobContract.get.queryOptions({ input: jobInput }),
initialData: initialJob,
refetchInterval: (query) => {
const status = query.state.data?.status
return status === 'queued' || status === 'running' ? POLL_INTERVAL : false
},
})
const retryMutation = useMutation({
...jobContract.post.mutationOptions(),
onMutate: () => setRetryError(undefined),
onSuccess: () => void refetch(),
onError: (error) => setRetryError(getErrorMessage(error)),
})
const completed = job.completed_documents + job.completed_sources
const total = job.total_documents + job.total_sources
const progress = total > 0 ? Math.round((completed / total) * 100) : 0
const statusLabel = {
queued: t(($) => $['newKnowledge.documentStatus.queued'], { ns: 'dataset' }),
running: t(($) => $['newKnowledge.documentStatus.processing'], { ns: 'dataset' }),
succeeded: t(($) => $['newKnowledge.processingTaskState.succeeded'], { ns: 'dataset' }),
failed: t(($) => $['newKnowledge.documentStatus.failed'], { ns: 'dataset' }),
}[job.status]
const errorMessage =
retryError ||
job.last_error_message ||
job.last_error_code ||
(statusError && getErrorMessage(statusError))
return (
<div className="flex flex-col gap-y-3 rounded-xl border border-components-panel-border bg-components-panel-bg p-4">
<div className="flex items-center justify-between gap-x-4">
<div className="text-sm font-medium text-text-primary">
{statusLabel}
<span className="ml-2 text-text-tertiary">{job.stage}</span>
</div>
<div className="text-xs text-text-secondary">{progress}%</div>
</div>
<div
className="h-1.5 overflow-hidden rounded-full bg-components-progress-bar-bg"
role="progressbar"
aria-valuemin={0}
aria-valuemax={100}
aria-valuenow={progress}
>
<div
className="h-full rounded-full bg-components-progress-bar-progress-solid"
style={{ width: `${progress}%` }}
/>
</div>
<div className="flex gap-x-5 text-xs text-text-secondary">
<span>
{t(($) => $['newKnowledge.documents'], { ns: 'dataset' })}: {job.completed_documents}/
{job.total_documents}
</span>
<span>
{t(($) => $['newKnowledge.sources'], { ns: 'dataset' })}: {job.completed_sources}/
{job.total_sources}
</span>
</div>
{errorMessage && (
<div className="text-xs break-words text-text-destructive">
{t(($) => $.error, { ns: 'common' })}: {errorMessage}
</div>
)}
{job.status === 'failed' && (
<div>
<Button
size="small"
variant="secondary"
loading={retryMutation.isPending}
onClick={() => retryMutation.mutate(jobInput)}
>
{t(($) => $.retry, { ns: 'dataset' })}
</Button>
</div>
)}
</div>
)
}
const KnowledgeFSUpgrade = ({ datasetId, disabled }: Props) => {
const { t } = useTranslation()
const [job, setJob] = useState<KnowledgeFsUpgradeJobResponse>()
const [startError, setStartError] = useState<string>()
const upgradeContract = consoleQuery.datasets.byDatasetId.knowledgeFsUpgrades
const upgradeInput = { params: { dataset_id: datasetId } }
const startMutation = useMutation({
...upgradeContract.post.mutationOptions(),
onMutate: () => setStartError(undefined),
onSuccess: setJob,
onError: (error) => setStartError(getErrorMessage(error)),
})
if (job) return <UpgradeProgress datasetId={datasetId} initialJob={job} />
return (
<div className="flex flex-col items-start gap-y-2">
<Button
variant="secondary"
loading={startMutation.isPending}
disabled={disabled || startMutation.isPending}
onClick={() => startMutation.mutate(upgradeInput)}
>
{t(($) => $['upgradeBtn.encourageShort'], { ns: 'billing' })}
</Button>
{startError && (
<div className="text-xs break-words text-text-destructive">
{t(($) => $.error, { ns: 'common' })}: {startError}
</div>
)}
</div>
)
}
export default KnowledgeFSUpgrade

View File

@ -5,6 +5,7 @@ import Divider from '@/app/components/base/divider'
import BasicInfoSection from './components/basic-info-section'
import ExternalKnowledgeSection from './components/external-knowledge-section'
import IndexingSection from './components/indexing-section'
import KnowledgeFSUpgrade from './components/knowledge-fs-upgrade'
import { useFormState } from './hooks/use-form-state'
const Form = () => {
@ -119,6 +120,12 @@ const Form = () => {
<Divider type="horizontal" className="my-1 h-px bg-divider-subtle" />
{!isExternalProvider && currentDataset && (
<KnowledgeFSUpgrade datasetId={currentDataset.id} disabled={readonly} />
)}
<Divider type="horizontal" className="my-1 h-px bg-divider-subtle" />
{/* Save Button */}
<div className="flex gap-x-1">
<div className="flex h-7 w-45 shrink-0 items-center pt-1" />