Refactor frontend architecture and simplify implementation

This commit is contained in:
Jyong 2026-07-29 02:30:27 -04:00
parent 5e7ea24340
commit 6cfd1abb59
19 changed files with 262 additions and 1358 deletions

View File

@ -33,24 +33,10 @@ class KnowledgeFSAccessDeniedHTTPError(BaseHTTPException):
code = 403
class KnowledgeFSRateLimitHTTPError(BaseHTTPException):
error_code = "knowledge_fs_rate_limit_exceeded"
description = "KnowledgeFS operation rate limit exceeded."
code = 429
class KnowledgeFSQuotaExceededHTTPError(BaseHTTPException):
error_code = "knowledge_fs_quota_exceeded"
description = "KnowledgeFS operation quota exceeded."
code = 403
__all__ = [
"KnowledgeFSAccessDeniedHTTPError",
"KnowledgeFSInvalidRequestHTTPError",
"KnowledgeFSOperationUnavailableHTTPError",
"KnowledgeFSQuotaExceededHTTPError",
"KnowledgeFSRateLimitHTTPError",
"KnowledgeFSSpaceNotFoundHTTPError",
"KnowledgeFSUpstreamUnavailableHTTPError",
]

View File

@ -24,12 +24,14 @@ from controllers.console.knowledge_fs.error import (
KnowledgeFSAccessDeniedHTTPError,
KnowledgeFSInvalidRequestHTTPError,
KnowledgeFSOperationUnavailableHTTPError,
KnowledgeFSQuotaExceededHTTPError,
KnowledgeFSRateLimitHTTPError,
KnowledgeFSSpaceNotFoundHTTPError,
KnowledgeFSUpstreamUnavailableHTTPError,
)
from controllers.console.wraps import account_initialization_required, setup_required
from controllers.console.wraps import (
account_initialization_required,
cloud_edition_billing_rate_limit_check,
setup_required,
)
from core.db.session_factory import session_factory
from libs.helper import dump_response
from libs.login import current_account_with_tenant, login_required
@ -44,10 +46,6 @@ from services.knowledge_fs.control_plane_service import (
from services.knowledge_fs.credential_service import (
KnowledgeFSCredentialPolicyError,
)
from services.knowledge_fs.operation_admission import (
KnowledgeFSOperationQuotaExceededError,
KnowledgeFSOperationRateLimitExceededError,
)
from services.knowledge_fs.product_authorization import (
KnowledgeFSProductNotFoundError,
)
@ -306,6 +304,7 @@ def _console_services() -> KnowledgeFSRuntime:
def _knowledge_fs_errors[**P, R](view: Callable[P, R]) -> Callable[P, R]:
@cloud_edition_billing_rate_limit_check("knowledge")
@wraps(view)
def decorated(*args: P.args, **kwargs: P.kwargs) -> R:
try:
@ -318,10 +317,6 @@ def _knowledge_fs_errors[**P, R](view: Callable[P, R]) -> Callable[P, R]:
raise NotFound() from exc
except KnowledgeFSProductRemoteError as exc:
raise KnowledgeFSUpstreamUnavailableHTTPError() from exc
except KnowledgeFSOperationRateLimitExceededError as exc:
raise KnowledgeFSRateLimitHTTPError() from exc
except KnowledgeFSOperationQuotaExceededError as exc:
raise KnowledgeFSQuotaExceededHTTPError() from exc
except KnowledgeFSProductRequestRejectedError as exc:
if exc.status_code == HTTPStatus.CONFLICT:
raise Conflict() from exc
@ -1978,7 +1973,7 @@ class KnowledgeFSSpaceQueryAdmissionApi(Resource):
raise KnowledgeFSOperationUnavailableError("KnowledgeFS direct query streaming is not configured")
actor_id, tenant_id = _actor()
payload = _payload(KnowledgeFSQueryCreatePayload)
issued = _console_services().direct_operation_admission.issue_interactive(
issued = _console_services().broker.issue_interactive(
tenant_id=tenant_id,
account_id=actor_id,
control_space_id=control_space_id,
@ -2473,7 +2468,7 @@ class KnowledgeFSSpaceUploadCapabilitiesApi(Resource):
raise KnowledgeFSOperationUnavailableError("KnowledgeFS direct upload is not configured")
actor_id, tenant_id = _actor()
payload = _payload(KnowledgeFSUploadCapabilityPayload)
issued = _console_services().direct_operation_admission.issue_interactive(
issued = _console_services().broker.issue_interactive(
tenant_id=tenant_id,
account_id=actor_id,
control_space_id=control_space_id,
@ -2536,7 +2531,7 @@ class KnowledgeFSSpaceQueryStreamCapabilityApi(Resource):
if direct_origin is None:
raise KnowledgeFSOperationUnavailableError("KnowledgeFS direct query streaming is not configured")
actor_id, tenant_id = _actor()
issued = _console_services().direct_operation_admission.issue_interactive(
issued = _console_services().broker.issue_interactive(
tenant_id=tenant_id,
account_id=actor_id,
control_space_id=control_space_id,
@ -2571,7 +2566,7 @@ class KnowledgeFSTaskStreamCapabilityApi(Resource):
raise KnowledgeFSOperationUnavailableError("KnowledgeFS direct streaming is not configured")
actor_id, tenant_id = _actor()
payload = _payload(KnowledgeFSStreamCapabilityPayload)
issued = _console_services().direct_operation_admission.issue_interactive(
issued = _console_services().broker.issue_interactive(
tenant_id=tenant_id,
account_id=actor_id,
control_space_id=payload.control_space_id,

View File

@ -27,23 +27,9 @@ class KnowledgeFSServiceInvalidRequestHTTPError(BaseHTTPException):
code = 400
class KnowledgeFSServiceRateLimitHTTPError(BaseHTTPException):
error_code = "knowledge_fs_rate_limit_exceeded"
description = "KnowledgeFS operation rate limit exceeded."
code = 429
class KnowledgeFSServiceQuotaExceededHTTPError(BaseHTTPException):
error_code = "knowledge_fs_quota_exceeded"
description = "KnowledgeFS operation quota exceeded."
code = 403
__all__ = [
"KnowledgeFSInvalidCredentialHTTPError",
"KnowledgeFSServiceInvalidRequestHTTPError",
"KnowledgeFSServiceOperationUnavailableHTTPError",
"KnowledgeFSServiceQuotaExceededHTTPError",
"KnowledgeFSServiceRateLimitHTTPError",
"KnowledgeFSServiceUpstreamUnavailableHTTPError",
]

View File

@ -23,8 +23,6 @@ from controllers.service_api.knowledge_fs.error import (
KnowledgeFSInvalidCredentialHTTPError,
KnowledgeFSServiceInvalidRequestHTTPError,
KnowledgeFSServiceOperationUnavailableHTTPError,
KnowledgeFSServiceQuotaExceededHTTPError,
KnowledgeFSServiceRateLimitHTTPError,
KnowledgeFSServiceUpstreamUnavailableHTTPError,
)
from core.db.session_factory import session_factory
@ -33,10 +31,6 @@ from services.knowledge_fs.credential_service import (
KnowledgeFSCredentialValidationError,
KnowledgeFSServiceCredentialProfile,
)
from services.knowledge_fs.operation_admission import (
KnowledgeFSOperationQuotaExceededError,
KnowledgeFSOperationRateLimitExceededError,
)
from services.knowledge_fs.product_dto import (
KnowledgeFSAdmittedQueryRequest,
KnowledgeFSAnswerTraceResponse,
@ -176,10 +170,6 @@ def _service_api_errors[**P, R](view: Callable[P, R]) -> Callable[P, R]:
raise NotFound() from exc
except KnowledgeFSProductRemoteError as exc:
raise KnowledgeFSServiceUpstreamUnavailableHTTPError() from exc
except KnowledgeFSOperationRateLimitExceededError as exc:
raise KnowledgeFSServiceRateLimitHTTPError() from exc
except KnowledgeFSOperationQuotaExceededError as exc:
raise KnowledgeFSServiceQuotaExceededHTTPError() from exc
except ValidationError as exc:
raise KnowledgeFSServiceInvalidRequestHTTPError() from exc
@ -551,7 +541,7 @@ class KnowledgeFSServiceQueryAdmissionApi(Resource):
runtime = _runtime()
profile = _profile(runtime, operation_id="createQuery", control_space_id=control_space_id)
payload = _payload(KnowledgeFSQueryCreatePayload)
issued = runtime.direct_operation_admission.issue_service(profile=profile, operation_id="createQuery")
issued = runtime.broker.issue_service(profile=profile, operation_id="createQuery")
admitted_request = KnowledgeFSAdmittedQueryRequest.model_validate(
{**payload.model_dump(mode="json", by_alias=True), "knowledgeSpaceId": issued.knowledge_space_id}
)

View File

@ -2,8 +2,6 @@
from __future__ import annotations
from collections.abc import Generator
from contextlib import contextmanager
from typing import Literal, cast
from pydantic import BaseModel, ConfigDict, Field, JsonValue, field_validator
@ -12,7 +10,6 @@ from core.app.entities.app_invoke_entities import DifyRunContext
from models.knowledge_fs import KnowledgeFSAppSpaceJoinType
from services.knowledge_fs.app_admission_service import KnowledgeFSAppAdmissionService
from services.knowledge_fs.capability_broker import KnowledgeFSCapabilityBroker, KnowledgeFSIssuedProductCapability
from services.knowledge_fs.operation_admission import KnowledgeFSOperationAdmissionService
from services.knowledge_fs.product_dto import (
KnowledgeFSResearchTaskCreatePayload,
KnowledgeFSResearchTaskResponse,
@ -48,12 +45,10 @@ class KnowledgeFSAppExecutionCapabilityService:
*,
admission: KnowledgeFSAppAdmissionService,
broker: KnowledgeFSCapabilityBroker,
operation_admission: KnowledgeFSOperationAdmissionService,
remote: KnowledgeFSProductRemotePort,
) -> None:
self._admission = admission
self._broker = broker
self._operation_admission = operation_admission
self._remote = remote
def issue(
@ -100,45 +95,32 @@ class KnowledgeFSAppExecutionCapabilityService:
or "{" in operation.kfs_path
):
raise KnowledgeFSOperationUnavailableError("KnowledgeFS app Research task creation is unavailable")
with self._admitted(tenant_id=run_context.tenant_id, operation_id=operation_id):
issued = self.issue(
tenant_id=run_context.tenant_id,
app_id=run_context.app_id,
control_space_id=resource.control_space_id,
caller_kind=caller_kind,
issued = self.issue(
tenant_id=run_context.tenant_id,
app_id=run_context.app_id,
control_space_id=resource.control_space_id,
caller_kind=caller_kind,
operation_id=operation_id,
trace_id=run_context.trace_session_id,
)
remote_payload = cast(
dict[str, JsonValue],
payload.model_dump(mode="json", exclude_none=True, by_alias=True),
)
remote_payload["knowledgeSpaceId"] = issued.knowledge_space_id
raw = self._remote.execute_json(
KnowledgeFSRemoteJSONRequest(
operation_id=operation_id,
trace_id=run_context.trace_session_id,
method=operation.method,
path=operation.kfs_path,
namespace_id=run_context.tenant_id,
knowledge_space_id=issued.knowledge_space_id,
capability_token=issued.token,
trace_id=issued.trace_id,
payload=remote_payload,
)
remote_payload = cast(
dict[str, JsonValue],
payload.model_dump(mode="json", exclude_none=True, by_alias=True),
)
remote_payload["knowledgeSpaceId"] = issued.knowledge_space_id
raw = self._remote.execute_json(
KnowledgeFSRemoteJSONRequest(
operation_id=operation_id,
method=operation.method,
path=operation.kfs_path,
namespace_id=run_context.tenant_id,
knowledge_space_id=issued.knowledge_space_id,
capability_token=issued.token,
trace_id=issued.trace_id,
payload=remote_payload,
)
)
response = KnowledgeFSResearchTaskResponse.model_validate(raw)
return response
@contextmanager
def _admitted(self, *, tenant_id: str, operation_id: str) -> Generator[None, None, None]:
charge = self._operation_admission.reserve(tenant_id=tenant_id, operation_id=operation_id)
try:
yield
except BaseException:
charge.refund()
raise
else:
charge.commit()
)
return KnowledgeFSResearchTaskResponse.model_validate(raw)
__all__ = ["KnowledgeFSAppExecutionCapabilityService", "KnowledgeResourceRef"]

View File

@ -2,15 +2,13 @@
from __future__ import annotations
from collections.abc import Callable, Generator
from contextlib import contextmanager
from collections.abc import Callable
from typing import Literal
from pydantic import BaseModel, JsonValue
from services.knowledge_fs.capability_broker import KnowledgeFSCapabilityBroker
from services.knowledge_fs.credential_service import KnowledgeFSServiceCredentialProfile
from services.knowledge_fs.operation_admission import KnowledgeFSOperationAdmissionService
from services.knowledge_fs.product_dto import (
KnowledgeFSAnswerTraceResponse,
KnowledgeFSBackgroundTaskListResponse,
@ -107,11 +105,9 @@ class KnowledgeFSDataFacade:
def __init__(
self,
*,
admission: KnowledgeFSOperationAdmissionService,
broker: KnowledgeFSCapabilityBroker,
remote: KnowledgeFSProductRemotePort,
) -> None:
self._admission = admission
self._broker = broker
self._remote = remote
@ -322,40 +318,39 @@ class KnowledgeFSDataFacade:
) -> KnowledgeFSDocumentUploadAcceptedResponse:
operation_id = "createDocument"
_assert_multipart_bff_ready(operation_id)
with self._admitted(tenant_id=tenant_id, operation_id=operation_id):
issued = self._broker.issue_interactive(
tenant_id=tenant_id,
account_id=account_id,
control_space_id=control_space_id,
issued = self._broker.issue_interactive(
tenant_id=tenant_id,
account_id=account_id,
control_space_id=control_space_id,
operation_id=operation_id,
)
operation = KNOWLEDGE_FS_PRODUCT_OPERATIONS[operation_id]
upload = body_reader(operation.max_request_bytes)
if not upload.filename or not upload.content_type or not isinstance(upload.body, bytes) or not upload.body:
raise KnowledgeFSProductRequestRejectedError(status_code=422)
if len(upload.body) > operation.max_request_bytes:
raise KnowledgeFSProductRequestRejectedError(status_code=413)
if operation.kfs_path is None:
raise KnowledgeFSOperationUnavailableError(f"KnowledgeFS operation is unavailable: {operation_id}")
path = _resolve_product_path(
template=operation.kfs_path,
knowledge_space_id=issued.knowledge_space_id,
resource_id=None,
resource_resolver=operation.resource_resolver,
path_parameters=(),
)
raw = self._remote.execute_multipart(
KnowledgeFSRemoteMultipartRequest(
operation_id=operation_id,
)
operation = KNOWLEDGE_FS_PRODUCT_OPERATIONS[operation_id]
upload = body_reader(operation.max_request_bytes)
if not upload.filename or not upload.content_type or not isinstance(upload.body, bytes) or not upload.body:
raise KnowledgeFSProductRequestRejectedError(status_code=422)
if len(upload.body) > operation.max_request_bytes:
raise KnowledgeFSProductRequestRejectedError(status_code=413)
if operation.kfs_path is None:
raise KnowledgeFSOperationUnavailableError(f"KnowledgeFS operation is unavailable: {operation_id}")
path = _resolve_product_path(
template=operation.kfs_path,
method=operation.method,
path=path,
namespace_id=tenant_id,
knowledge_space_id=issued.knowledge_space_id,
resource_id=None,
resource_resolver=operation.resource_resolver,
path_parameters=(),
)
raw = self._remote.execute_multipart(
KnowledgeFSRemoteMultipartRequest(
operation_id=operation_id,
method=operation.method,
path=path,
namespace_id=tenant_id,
knowledge_space_id=issued.knowledge_space_id,
capability_token=issued.token,
trace_id=issued.trace_id,
file=upload,
)
capability_token=issued.token,
trace_id=issued.trace_id,
file=upload,
)
)
return KnowledgeFSDocumentUploadAcceptedResponse.model_validate(raw)
def upload_small_file(
@ -369,42 +364,41 @@ class KnowledgeFSDataFacade:
) -> KnowledgeFSSmallFileUploadResponse:
operation_id = "uploadSmallFile"
_assert_binary_bff_ready(operation_id)
with self._admitted(tenant_id=tenant_id, operation_id=operation_id):
issued = self._broker.issue_interactive(
tenant_id=tenant_id,
account_id=account_id,
control_space_id=control_space_id,
issued = self._broker.issue_interactive(
tenant_id=tenant_id,
account_id=account_id,
control_space_id=control_space_id,
operation_id=operation_id,
resource_id=upload_session_id,
)
operation = KNOWLEDGE_FS_PRODUCT_OPERATIONS[operation_id]
body = body_reader(operation.max_request_bytes)
if not isinstance(body, bytes) or not body:
raise KnowledgeFSProductRequestRejectedError(status_code=422)
if len(body) > operation.max_request_bytes:
raise KnowledgeFSProductRequestRejectedError(status_code=413)
if operation.kfs_path is None:
raise KnowledgeFSOperationUnavailableError(f"KnowledgeFS operation is unavailable: {operation_id}")
path = _resolve_product_path(
template=operation.kfs_path,
knowledge_space_id=issued.knowledge_space_id,
resource_id=upload_session_id,
resource_resolver=operation.resource_resolver,
path_parameters=(),
)
raw = self._remote.execute_binary(
KnowledgeFSRemoteBinaryRequest(
operation_id=operation_id,
resource_id=upload_session_id,
)
operation = KNOWLEDGE_FS_PRODUCT_OPERATIONS[operation_id]
body = body_reader(operation.max_request_bytes)
if not isinstance(body, bytes) or not body:
raise KnowledgeFSProductRequestRejectedError(status_code=422)
if len(body) > operation.max_request_bytes:
raise KnowledgeFSProductRequestRejectedError(status_code=413)
if operation.kfs_path is None:
raise KnowledgeFSOperationUnavailableError(f"KnowledgeFS operation is unavailable: {operation_id}")
path = _resolve_product_path(
template=operation.kfs_path,
method=operation.method,
path=path,
namespace_id=tenant_id,
knowledge_space_id=issued.knowledge_space_id,
resource_id=upload_session_id,
resource_resolver=operation.resource_resolver,
path_parameters=(),
)
raw = self._remote.execute_binary(
KnowledgeFSRemoteBinaryRequest(
operation_id=operation_id,
method=operation.method,
path=path,
namespace_id=tenant_id,
knowledge_space_id=issued.knowledge_space_id,
capability_token=issued.token,
trace_id=issued.trace_id,
body=body,
query=(("knowledgeSpaceId", issued.knowledge_space_id),),
)
capability_token=issued.token,
trace_id=issued.trace_id,
body=body,
query=(("knowledgeSpaceId", issued.knowledge_space_id),),
)
)
return KnowledgeFSSmallFileUploadResponse.model_validate(raw)
def get_document(
@ -1509,22 +1503,21 @@ class KnowledgeFSDataFacade:
headers: tuple[tuple[str, str], ...] = (),
) -> JsonValue:
_assert_json_bff_ready(operation_id)
with self._admitted(tenant_id=profile.tenant_id, operation_id=operation_id):
issued = self._broker.issue_service(profile=profile, operation_id=operation_id, resource_id=resource_id)
return self._execute(
operation_id=operation_id,
namespace_id=profile.tenant_id,
knowledge_space_id=profile.knowledge_space_id,
knowledge_space_revision=profile.knowledge_space_revision,
capability_token=issued.token,
trace_id=issued.trace_id,
payload=payload,
query=query,
bind_space_in_body=bind_space_in_body,
resource_id=resource_id,
path_parameters=path_parameters,
headers=headers,
)
issued = self._broker.issue_service(profile=profile, operation_id=operation_id, resource_id=resource_id)
return self._execute(
operation_id=operation_id,
namespace_id=profile.tenant_id,
knowledge_space_id=profile.knowledge_space_id,
knowledge_space_revision=profile.knowledge_space_revision,
capability_token=issued.token,
trace_id=issued.trace_id,
payload=payload,
query=query,
bind_space_in_body=bind_space_in_body,
resource_id=resource_id,
path_parameters=path_parameters,
headers=headers,
)
def _interactive(
self,
@ -1541,39 +1534,27 @@ class KnowledgeFSDataFacade:
headers: tuple[tuple[str, str], ...] = (),
) -> JsonValue:
_assert_json_bff_ready(operation_id)
with self._admitted(tenant_id=tenant_id, operation_id=operation_id):
issued = self._broker.issue_interactive(
tenant_id=tenant_id,
account_id=account_id,
control_space_id=control_space_id,
operation_id=operation_id,
resource_id=resource_id,
)
return self._execute(
operation_id=operation_id,
namespace_id=tenant_id,
knowledge_space_id=issued.knowledge_space_id,
knowledge_space_revision=issued.knowledge_space_revision,
capability_token=issued.token,
trace_id=issued.trace_id,
payload=payload,
query=query,
bind_space_in_body=bind_space_in_body,
resource_id=resource_id,
path_parameters=path_parameters,
headers=headers,
)
@contextmanager
def _admitted(self, *, tenant_id: str, operation_id: str) -> Generator[None]:
charge = self._admission.reserve(tenant_id=tenant_id, operation_id=operation_id)
try:
yield
except BaseException:
charge.refund()
raise
else:
charge.commit()
issued = self._broker.issue_interactive(
tenant_id=tenant_id,
account_id=account_id,
control_space_id=control_space_id,
operation_id=operation_id,
resource_id=resource_id,
)
return self._execute(
operation_id=operation_id,
namespace_id=tenant_id,
knowledge_space_id=issued.knowledge_space_id,
knowledge_space_revision=issued.knowledge_space_revision,
capability_token=issued.token,
trace_id=issued.trace_id,
payload=payload,
query=query,
bind_space_in_body=bind_space_in_body,
resource_id=resource_id,
path_parameters=path_parameters,
headers=headers,
)
def _interactive_child(
self,

View File

@ -40,13 +40,6 @@ class KnowledgeFSLifecycleTaskMetric(NamedTuple):
status: Literal["dispatch_error", "queued", "retry", "running", "succeeded"]
class KnowledgeFSOperationAdmissionMetric(NamedTuple):
operation_id: str
bucket: str
phase: Literal["commit", "refund", "reserve"]
outcome: Literal["failure", "success"]
class KnowledgeFSOperationalMetricsPort(Protocol):
def record_batch_status(self, event: KnowledgeFSBatchStatusMetric) -> None: ...
@ -56,8 +49,6 @@ class KnowledgeFSOperationalMetricsPort(Protocol):
def record_lifecycle_task(self, event: KnowledgeFSLifecycleTaskMetric) -> None: ...
def record_operation_admission(self, event: KnowledgeFSOperationAdmissionMetric) -> None: ...
def register_control_space_state_gauge(self, read_counts: Callable[[], Mapping[str, int]]) -> None: ...
@ -136,11 +127,6 @@ class OpenTelemetryKnowledgeFSOperationalMetrics:
description="KnowledgeFS capability revoke enqueue-to-ack latency",
unit="s",
)
self._operation_admission = resolved_meter.create_counter(
"dify.knowledge_fs.operation_admission",
description="KnowledgeFS direct-operation reserve and finalization outcomes",
unit="{operation}",
)
def record_capability_issuance(self, event: KnowledgeFSCapabilityIssuanceMetric) -> None:
self._capability_issuance.add(
@ -174,17 +160,6 @@ class OpenTelemetryKnowledgeFSOperationalMetrics:
if event.operation == "revoke" and event.status == "succeeded":
self._revoke_latency.record(event.duration_seconds, attributes={"operation": "revoke"})
def record_operation_admission(self, event: KnowledgeFSOperationAdmissionMetric) -> None:
self._operation_admission.add(
1,
attributes={
"bucket": event.bucket,
"operation_id": event.operation_id,
"outcome": event.outcome,
"phase": event.phase,
},
)
def register_control_space_state_gauge(self, read_counts: Callable[[], Mapping[str, int]]) -> None:
"""Register one DB-backed current-state instrument per process."""
@ -227,7 +202,6 @@ __all__ = [
"KnowledgeFSCapabilityIssuanceMetric",
"KnowledgeFSControlSpaceStateMetric",
"KnowledgeFSLifecycleTaskMetric",
"KnowledgeFSOperationAdmissionMetric",
"KnowledgeFSOperationalMetricsPort",
"OpenTelemetryKnowledgeFSOperationalMetrics",
"get_knowledge_fs_operational_metrics",

View File

@ -1,408 +0,0 @@
"""Weighted admission for BFF calls and browser-executed KnowledgeFS operations."""
from __future__ import annotations
import logging
import time
import uuid
from collections.abc import Callable
from typing import Literal, NamedTuple, Protocol, cast
from configs import dify_config
from extensions.ext_redis import redis_client
from services.billing_service import BillingService
from services.feature_service import FeatureService, KnowledgeRateLimitModel
from services.knowledge_fs.capability_broker import KnowledgeFSIssuedProductCapability
from services.knowledge_fs.credential_service import KnowledgeFSServiceCredentialProfile
from services.knowledge_fs.observability import (
KnowledgeFSOperationAdmissionMetric,
KnowledgeFSOperationalMetricsPort,
get_knowledge_fs_operational_metrics,
)
from services.knowledge_fs.product_operations import KNOWLEDGE_FS_PRODUCT_OPERATIONS
logger = logging.getLogger(__name__)
_BILLING_FEATURE_KEY = "knowledge_fs_operations"
_RATE_LIMIT_WINDOW_MS = 60_000
_WEIGHTED_RATE_LIMIT_SCRIPT = """
redis.call('ZREMRANGEBYSCORE', KEYS[1], 0, ARGV[2])
local current = redis.call('ZCARD', KEYS[1])
local requested = tonumber(ARGV[4])
if current + requested > tonumber(ARGV[3]) then
return 0
end
for index = 1, requested do
redis.call('ZADD', KEYS[1], ARGV[1], ARGV[5] .. ':' .. index)
end
redis.call('PEXPIRE', KEYS[1], 61000)
return 1
"""
class KnowledgeFSOperationRateLimitExceededError(RuntimeError):
"""The operation's weighted per-minute bucket is exhausted."""
class KnowledgeFSOperationQuotaExceededError(RuntimeError):
"""The billing service explicitly rejected an operation reservation."""
class KnowledgeFSOperationUsage(NamedTuple):
tenant_id: str
operation_id: str
bucket: str
billing_cost: int
rate_limit_cost: int
class KnowledgeFSOperationChargePort(Protocol):
def commit(self) -> None: ...
def refund(self) -> None: ...
class KnowledgeFSOperationRateLimitPort(Protocol):
def admit(self, usage: KnowledgeFSOperationUsage) -> None: ...
class KnowledgeFSOperationBillingPort(Protocol):
def reserve(self, usage: KnowledgeFSOperationUsage) -> KnowledgeFSOperationChargePort: ...
class KnowledgeFSRateLimitAuditPort(Protocol):
def record_rejection(self, usage: KnowledgeFSOperationUsage, *, subscription_plan: str) -> None: ...
class KnowledgeFSBillingGateway(Protocol):
def quota_reserve(self, **kwargs: object) -> dict[str, object]: ...
def quota_commit(self, **kwargs: object) -> dict[str, object]: ...
def quota_release(self, **kwargs: object) -> dict[str, object]: ...
class KnowledgeFSRedisEvalPort(Protocol):
def eval(self, *args: object) -> int: ...
class _NoopCharge:
def commit(self) -> None:
return
def refund(self) -> None:
return
class _DifyBillingCharge:
def __init__(
self,
*,
gateway: KnowledgeFSBillingGateway,
usage: KnowledgeFSOperationUsage,
reservation_id: str,
) -> None:
self._gateway = gateway
self._usage = usage
self._reservation_id = reservation_id
self._finalized = False
def commit(self) -> None:
if self._finalized:
return
self._finalized = True
try:
self._gateway.quota_commit(
tenant_id=self._usage.tenant_id,
feature_key=_BILLING_FEATURE_KEY,
bucket=self._usage.bucket,
reservation_id=self._reservation_id,
actual_amount=self._usage.billing_cost,
meta={"operation_id": self._usage.operation_id, "source": "knowledge_fs"},
)
except Exception:
logger.exception(
"KnowledgeFS billing commit failed for tenant_id=%s operation_id=%s",
self._usage.tenant_id,
self._usage.operation_id,
)
def refund(self) -> None:
if self._finalized:
return
self._finalized = True
try:
self._gateway.quota_release(
tenant_id=self._usage.tenant_id,
feature_key=_BILLING_FEATURE_KEY,
bucket=self._usage.bucket,
reservation_id=self._reservation_id,
)
except Exception:
logger.exception(
"KnowledgeFS billing release failed for tenant_id=%s operation_id=%s",
self._usage.tenant_id,
self._usage.operation_id,
)
class LoggingKnowledgeFSRateLimitAudit:
def record_rejection(self, usage: KnowledgeFSOperationUsage, *, subscription_plan: str) -> None:
logger.warning(
"KnowledgeFS weighted rate limit rejected tenant_id=%s operation_id=%s bucket=%s cost=%s plan=%s",
usage.tenant_id,
usage.operation_id,
usage.bucket,
usage.rate_limit_cost,
subscription_plan,
)
class DifyKnowledgeFSWeightedRateLimitPort:
"""Consume weighted operation units atomically from one Redis bucket."""
def __init__(
self,
*,
redis: KnowledgeFSRedisEvalPort | None = None,
audit: KnowledgeFSRateLimitAuditPort,
rate_limit_lookup: Callable[[str], KnowledgeRateLimitModel] = FeatureService.get_knowledge_rate_limit,
clock_ms: Callable[[], int] = lambda: int(time.time() * 1000),
member_id: Callable[[], str] = lambda: str(uuid.uuid4()),
) -> None:
self._redis = redis or cast(KnowledgeFSRedisEvalPort, redis_client)
self._audit = audit
self._rate_limit_lookup = rate_limit_lookup
self._clock_ms = clock_ms
self._member_id = member_id
def admit(self, usage: KnowledgeFSOperationUsage) -> None:
rate_limit = self._rate_limit_lookup(usage.tenant_id)
if not rate_limit.enabled:
return
now = self._clock_ms()
accepted = self._redis.eval(
_WEIGHTED_RATE_LIMIT_SCRIPT,
1,
f"knowledge_fs:rate_limit:{usage.tenant_id}:{usage.bucket}",
now,
now - _RATE_LIMIT_WINDOW_MS,
rate_limit.limit,
usage.rate_limit_cost,
self._member_id(),
)
if accepted == 1:
return
self._audit.record_rejection(usage, subscription_plan=rate_limit.subscription_plan)
raise KnowledgeFSOperationRateLimitExceededError("KnowledgeFS operation rate limit exceeded")
class DifyKnowledgeFSBillingPort:
"""Reserve operation-specific billing units for the caller to commit or release."""
def __init__(
self,
*,
gateway: KnowledgeFSBillingGateway | None = None,
billing_enabled: Callable[[], bool] = lambda: dify_config.BILLING_ENABLED,
request_id: Callable[[], str] = lambda: str(uuid.uuid4()),
) -> None:
self._gateway = gateway or cast(KnowledgeFSBillingGateway, BillingService)
self._billing_enabled = billing_enabled
self._request_id = request_id
def reserve(self, usage: KnowledgeFSOperationUsage) -> KnowledgeFSOperationChargePort:
if not self._billing_enabled():
return _NoopCharge()
try:
result = self._gateway.quota_reserve(
tenant_id=usage.tenant_id,
feature_key=_BILLING_FEATURE_KEY,
bucket=usage.bucket,
request_id=self._request_id(),
amount=usage.billing_cost,
meta={"operation_id": usage.operation_id, "source": "knowledge_fs"},
)
except Exception:
logger.exception(
"KnowledgeFS billing reservation unavailable for tenant_id=%s operation_id=%s; allowing request",
usage.tenant_id,
usage.operation_id,
)
return _NoopCharge()
reservation_id = result.get("reservation_id")
if not isinstance(reservation_id, str) or not reservation_id:
raise KnowledgeFSOperationQuotaExceededError("KnowledgeFS operation quota exceeded")
return _DifyBillingCharge(gateway=self._gateway, usage=usage, reservation_id=reservation_id)
class KnowledgeFSOperationAdmissionService:
def __init__(
self,
*,
rate_limit: KnowledgeFSOperationRateLimitPort,
billing: KnowledgeFSOperationBillingPort,
) -> None:
self._rate_limit = rate_limit
self._billing = billing
def reserve(self, *, tenant_id: str, operation_id: str) -> KnowledgeFSOperationChargePort:
operation = KNOWLEDGE_FS_PRODUCT_OPERATIONS[operation_id]
usage = KnowledgeFSOperationUsage(
tenant_id=tenant_id,
operation_id=operation_id,
bucket=operation.rate_limit_bucket,
billing_cost=operation.billing_cost,
rate_limit_cost=operation.rate_limit_cost,
)
self._rate_limit.admit(usage)
return self._billing.reserve(usage)
class KnowledgeFSDirectCapabilityBrokerPort(Protocol):
def issue_interactive(
self,
*,
tenant_id: str,
account_id: str,
control_space_id: str,
operation_id: str,
resource_id: str | None = None,
trace_id: str | None = None,
) -> KnowledgeFSIssuedProductCapability: ...
def issue_service(
self,
*,
profile: KnowledgeFSServiceCredentialProfile,
operation_id: str,
resource_id: str | None = None,
trace_id: str | None = None,
) -> KnowledgeFSIssuedProductCapability: ...
class KnowledgeFSDirectOperationAdmissionService:
"""Admit browser-executed operations at the Dify Capability issuance boundary.
A successful Capability issuance is the only outcome Dify can observe before the browser
talks to KnowledgeFS, so it commits the operation charge immediately. Issuance failures
release the reservation. Controllers must treat alternative endpoints for the same manifest
operation as separate entry points and must never chain them for one user action.
"""
_admission: KnowledgeFSOperationAdmissionService
_broker: KnowledgeFSDirectCapabilityBrokerPort
_metrics: KnowledgeFSOperationalMetricsPort
def __init__(
self,
*,
admission: KnowledgeFSOperationAdmissionService,
broker: KnowledgeFSDirectCapabilityBrokerPort,
metrics: KnowledgeFSOperationalMetricsPort | None = None,
) -> None:
self._admission = admission
self._broker = broker
self._metrics = metrics or get_knowledge_fs_operational_metrics()
def issue_interactive(
self,
*,
tenant_id: str,
account_id: str,
control_space_id: str,
operation_id: str,
resource_id: str | None = None,
trace_id: str | None = None,
) -> KnowledgeFSIssuedProductCapability:
return self._issue(
tenant_id=tenant_id,
operation_id=operation_id,
issue=lambda: self._broker.issue_interactive(
tenant_id=tenant_id,
account_id=account_id,
control_space_id=control_space_id,
operation_id=operation_id,
resource_id=resource_id,
trace_id=trace_id,
),
)
def issue_service(
self,
*,
profile: KnowledgeFSServiceCredentialProfile,
operation_id: str,
resource_id: str | None = None,
trace_id: str | None = None,
) -> KnowledgeFSIssuedProductCapability:
return self._issue(
tenant_id=profile.tenant_id,
operation_id=operation_id,
issue=lambda: self._broker.issue_service(
profile=profile,
operation_id=operation_id,
resource_id=resource_id,
trace_id=trace_id,
),
)
def _issue(
self,
*,
tenant_id: str,
operation_id: str,
issue: Callable[[], KnowledgeFSIssuedProductCapability],
) -> KnowledgeFSIssuedProductCapability:
bucket = KNOWLEDGE_FS_PRODUCT_OPERATIONS[operation_id].rate_limit_bucket
try:
charge = self._admission.reserve(tenant_id=tenant_id, operation_id=operation_id)
except BaseException:
self._record_metric(operation_id, bucket=bucket, phase="reserve", outcome="failure")
raise
self._record_metric(operation_id, bucket=bucket, phase="reserve", outcome="success")
try:
issued = issue()
except BaseException:
try:
charge.refund()
except BaseException:
self._record_metric(operation_id, bucket=bucket, phase="refund", outcome="failure")
raise
self._record_metric(operation_id, bucket=bucket, phase="refund", outcome="success")
raise
try:
charge.commit()
except BaseException:
self._record_metric(operation_id, bucket=bucket, phase="commit", outcome="failure")
raise
self._record_metric(operation_id, bucket=bucket, phase="commit", outcome="success")
return issued
def _record_metric(
self,
operation_id: str,
*,
bucket: str,
phase: Literal["commit", "refund", "reserve"],
outcome: Literal["failure", "success"],
) -> None:
try:
self._metrics.record_operation_admission(
KnowledgeFSOperationAdmissionMetric(operation_id, bucket, phase, outcome)
)
except Exception:
logger.warning("KnowledgeFS direct-operation admission metric export failed", exc_info=True)
__all__ = [
"DifyKnowledgeFSBillingPort",
"DifyKnowledgeFSWeightedRateLimitPort",
"KnowledgeFSDirectOperationAdmissionService",
"KnowledgeFSOperationAdmissionService",
"KnowledgeFSOperationChargePort",
"KnowledgeFSOperationQuotaExceededError",
"KnowledgeFSOperationRateLimitExceededError",
"KnowledgeFSOperationUsage",
"LoggingKnowledgeFSRateLimitAudit",
]

View File

@ -36,12 +36,9 @@ class KnowledgeFSProductOperation(NamedTuple):
"document", "job", "knowledge_space", "namespace", "query", "research_task", "source", "upload_session"
]
rbac_permission: KnowledgeFSProductPermission
billing_cost: int
max_request_bytes: int
max_response_bytes: int
stream_kind: Literal["buffered-multipart", "direct-upload", "json", "sse"]
rate_limit_bucket: Literal["direct", "import", "job", "query", "read", "write"]
rate_limit_cost: int
@property
def action(self) -> str | None:
@ -61,19 +58,10 @@ def _operation(
resource_resolver: Literal[
"document", "job", "knowledge_space", "namespace", "query", "research_task", "source", "upload_session"
],
billing_cost: int,
max_request_bytes: int,
max_response_bytes: int,
stream_kind: Literal["buffered-multipart", "direct-upload", "json", "sse"],
rate_limit_bucket: Literal["direct", "import", "job", "query", "read", "write"] | None = None,
rate_limit_cost: int | None = None,
) -> KnowledgeFSProductOperation:
resolved_rate_limit_bucket = rate_limit_bucket or _default_rate_limit_bucket(
method=method,
permission=permission,
resource_resolver=resource_resolver,
stream_kind=stream_kind,
)
return KnowledgeFSProductOperation(
method=method,
capability_operation_id=capability_operation_id,
@ -82,35 +70,12 @@ def _operation(
transport=transport,
resource_resolver=resource_resolver,
rbac_permission=permission,
billing_cost=billing_cost,
max_request_bytes=max_request_bytes,
max_response_bytes=max_response_bytes,
stream_kind=stream_kind,
rate_limit_bucket=resolved_rate_limit_bucket,
rate_limit_cost=rate_limit_cost if rate_limit_cost is not None else billing_cost,
)
def _default_rate_limit_bucket(
*,
method: Literal["DELETE", "GET", "PATCH", "POST", "PUT"],
permission: KnowledgeFSProductPermission,
resource_resolver: Literal[
"document", "job", "knowledge_space", "namespace", "query", "research_task", "source", "upload_session"
],
stream_kind: Literal["buffered-multipart", "direct-upload", "json", "sse"],
) -> Literal["direct", "import", "job", "query", "read", "write"]:
if stream_kind in {"direct-upload", "sse"}:
return "direct"
if permission == KnowledgeFSProductPermission.QUERY:
return "query"
if resource_resolver == "job":
return "job"
if method == "GET":
return "read"
return "write"
KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductOperation]] = MappingProxyType(
{
"batchSpaceSummaries": _operation(
@ -120,7 +85,6 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/internal/knowledge-spaces/product-summaries/batch",
"json",
resource_resolver="namespace",
billing_cost=2,
max_request_bytes=64 * 1024,
max_response_bytes=1024 * 1024,
stream_kind="json",
@ -132,7 +96,6 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/knowledge-spaces/{id}",
"json",
resource_resolver="knowledge_space",
billing_cost=1,
max_request_bytes=0,
max_response_bytes=256 * 1024,
stream_kind="json",
@ -144,7 +107,6 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/knowledge-spaces/{id}",
"json",
resource_resolver="knowledge_space",
billing_cost=3,
max_request_bytes=32 * 1024,
max_response_bytes=256 * 1024,
stream_kind="json",
@ -156,7 +118,6 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/knowledge-spaces/{id}/product-settings",
"json",
resource_resolver="knowledge_space",
billing_cost=1,
max_request_bytes=0,
max_response_bytes=256 * 1024,
stream_kind="json",
@ -168,7 +129,6 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/knowledge-spaces/{id}/overview/stats",
"json",
resource_resolver="knowledge_space",
billing_cost=2,
max_request_bytes=16 * 1024,
max_response_bytes=256 * 1024,
stream_kind="json",
@ -180,7 +140,6 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/knowledge-spaces/{id}/overview/query-outcomes",
"json",
resource_resolver="knowledge_space",
billing_cost=3,
max_request_bytes=16 * 1024,
max_response_bytes=512 * 1024,
stream_kind="json",
@ -192,7 +151,6 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/knowledge-spaces/{id}/overview/inventory",
"json",
resource_resolver="knowledge_space",
billing_cost=3,
max_request_bytes=0,
max_response_bytes=256 * 1024,
stream_kind="json",
@ -204,7 +162,6 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/knowledge-spaces/{id}/overview/health",
"json",
resource_resolver="knowledge_space",
billing_cost=2,
max_request_bytes=0,
max_response_bytes=256 * 1024,
stream_kind="json",
@ -216,7 +173,6 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/knowledge-spaces/{id}/product-settings",
"json",
resource_resolver="knowledge_space",
billing_cost=5,
max_request_bytes=64 * 1024,
max_response_bytes=256 * 1024,
stream_kind="json",
@ -228,7 +184,6 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/knowledge-spaces/{id}/embedding-profile",
"json",
resource_resolver="knowledge_space",
billing_cost=5,
max_request_bytes=32 * 1024,
max_response_bytes=256 * 1024,
stream_kind="json",
@ -240,7 +195,6 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/knowledge-spaces/{id}/retrieval-profile",
"json",
resource_resolver="knowledge_space",
billing_cost=5,
max_request_bytes=64 * 1024,
max_response_bytes=256 * 1024,
stream_kind="json",
@ -252,7 +206,6 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/knowledge-spaces/{id}/profile-migrations/{migrationId}",
"json",
resource_resolver="knowledge_space",
billing_cost=1,
max_request_bytes=0,
max_response_bytes=256 * 1024,
stream_kind="json",
@ -264,7 +217,6 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/knowledge-spaces/{id}/documents",
"json",
resource_resolver="knowledge_space",
billing_cost=2,
max_request_bytes=16 * 1024,
max_response_bytes=2 * 1024 * 1024,
stream_kind="json",
@ -276,7 +228,6 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/knowledge-spaces/{id}/logical-documents",
"json",
resource_resolver="knowledge_space",
billing_cost=2,
max_request_bytes=16 * 1024,
max_response_bytes=2 * 1024 * 1024,
stream_kind="json",
@ -288,7 +239,6 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/knowledge-spaces/{id}/logical-documents/{documentId}",
"json",
resource_resolver="document",
billing_cost=2,
max_request_bytes=0,
max_response_bytes=512 * 1024,
stream_kind="json",
@ -300,7 +250,6 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/knowledge-spaces/{id}/documents",
"multipart",
resource_resolver="knowledge_space",
billing_cost=10,
max_request_bytes=15 * 1024 * 1024,
max_response_bytes=1024 * 1024,
stream_kind="buffered-multipart",
@ -312,7 +261,6 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/knowledge-spaces/{id}/documents/{documentId}",
"json",
resource_resolver="document",
billing_cost=2,
max_request_bytes=0,
max_response_bytes=512 * 1024,
stream_kind="json",
@ -324,7 +272,6 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/knowledge-spaces/{id}/documents/{documentId}/outline",
"json",
resource_resolver="document",
billing_cost=3,
max_request_bytes=0,
max_response_bytes=4 * 1024 * 1024,
stream_kind="json",
@ -336,7 +283,6 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/knowledge-spaces/{id}/documents/{documentId}/revisions",
"json",
resource_resolver="document",
billing_cost=2,
max_request_bytes=16 * 1024,
max_response_bytes=2 * 1024 * 1024,
stream_kind="json",
@ -348,7 +294,6 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/knowledge-spaces/{id}/documents/{documentId}/metadata",
"json",
resource_resolver="document",
billing_cost=4,
max_request_bytes=128 * 1024,
max_response_bytes=512 * 1024,
stream_kind="json",
@ -360,7 +305,6 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/knowledge-spaces/{id}/documents/{documentId}/revisions/{revision}/chunks",
"json",
resource_resolver="document",
billing_cost=3,
max_request_bytes=16 * 1024,
max_response_bytes=4 * 1024 * 1024,
stream_kind="json",
@ -372,7 +316,6 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/knowledge-spaces/{id}/documents/{documentId}/revisions/{revision}/chunks/{chunkId}",
"json",
resource_resolver="document",
billing_cost=2,
max_request_bytes=0,
max_response_bytes=1024 * 1024,
stream_kind="json",
@ -384,7 +327,6 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/knowledge-spaces/{id}/documents/{documentId}",
"json",
resource_resolver="document",
billing_cost=8,
max_request_bytes=32 * 1024,
max_response_bytes=512 * 1024,
stream_kind="json",
@ -396,7 +338,6 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/knowledge-spaces/{id}/documents/bulk",
"json",
resource_resolver="knowledge_space",
billing_cost=20,
max_request_bytes=1024 * 1024,
max_response_bytes=4 * 1024 * 1024,
stream_kind="json",
@ -408,11 +349,9 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/knowledge-spaces/{id}/documents/bulk/reindex",
"json",
resource_resolver="knowledge_space",
billing_cost=20,
max_request_bytes=1024 * 1024,
max_response_bytes=4 * 1024 * 1024,
stream_kind="json",
rate_limit_bucket="import",
),
"getCompilationJob": _operation(
"GET",
@ -421,7 +360,6 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/jobs/{id}",
"json",
resource_resolver="job",
billing_cost=1,
max_request_bytes=16 * 1024,
max_response_bytes=512 * 1024,
stream_kind="json",
@ -433,7 +371,6 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/jobs/{id}",
"json",
resource_resolver="job",
billing_cost=5,
max_request_bytes=16 * 1024,
max_response_bytes=512 * 1024,
stream_kind="json",
@ -445,7 +382,6 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/jobs/{id}/retry",
"json",
resource_resolver="job",
billing_cost=8,
max_request_bytes=16 * 1024,
max_response_bytes=512 * 1024,
stream_kind="json",
@ -457,7 +393,6 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/bulk-jobs/{id}",
"json",
resource_resolver="job",
billing_cost=1,
max_request_bytes=16 * 1024,
max_response_bytes=2 * 1024 * 1024,
stream_kind="json",
@ -469,7 +404,6 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/knowledge-spaces/{id}/background-tasks",
"json",
resource_resolver="knowledge_space",
billing_cost=2,
max_request_bytes=16 * 1024,
max_response_bytes=4 * 1024 * 1024,
stream_kind="json",
@ -481,7 +415,6 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/knowledge-spaces/{id}/background-tasks/{taskKind}/{taskId}/cancel",
"json",
resource_resolver="job",
billing_cost=5,
max_request_bytes=16 * 1024,
max_response_bytes=512 * 1024,
stream_kind="json",
@ -493,7 +426,6 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/knowledge-spaces/{id}/background-tasks/{taskKind}/{taskId}/retry",
"json",
resource_resolver="job",
billing_cost=8,
max_request_bytes=16 * 1024,
max_response_bytes=512 * 1024,
stream_kind="json",
@ -505,7 +437,6 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/knowledge-spaces/{id}/sources",
"json",
resource_resolver="knowledge_space",
billing_cost=2,
max_request_bytes=16 * 1024,
max_response_bytes=2 * 1024 * 1024,
stream_kind="json",
@ -517,7 +448,6 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/knowledge-spaces/{id}/sources",
"json",
resource_resolver="knowledge_space",
billing_cost=8,
max_request_bytes=256 * 1024,
max_response_bytes=512 * 1024,
stream_kind="json",
@ -529,7 +459,6 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/knowledge-spaces/{id}/sources/{sourceId}",
"json",
resource_resolver="source",
billing_cost=1,
max_request_bytes=0,
max_response_bytes=512 * 1024,
stream_kind="json",
@ -541,7 +470,6 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/knowledge-spaces/{id}/sources/{sourceId}",
"json",
resource_resolver="source",
billing_cost=4,
max_request_bytes=256 * 1024,
max_response_bytes=512 * 1024,
stream_kind="json",
@ -553,7 +481,6 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/knowledge-spaces/{id}/sources/{sourceId}",
"json",
resource_resolver="source",
billing_cost=8,
max_request_bytes=32 * 1024,
max_response_bytes=512 * 1024,
stream_kind="json",
@ -565,7 +492,6 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/knowledge-spaces/{id}/sources/{sourceId}/test",
"json",
resource_resolver="source",
billing_cost=3,
max_request_bytes=0,
max_response_bytes=256 * 1024,
stream_kind="json",
@ -577,11 +503,9 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/knowledge-spaces/{id}/sources/{sourceId}/sync",
"json",
resource_resolver="source",
billing_cost=8,
max_request_bytes=16 * 1024,
max_response_bytes=512 * 1024,
stream_kind="json",
rate_limit_bucket="import",
),
"listSourceProviders": _operation(
"GET",
@ -590,7 +514,6 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/source-providers",
"json",
resource_resolver="namespace",
billing_cost=1,
max_request_bytes=0,
max_response_bytes=512 * 1024,
stream_kind="json",
@ -602,7 +525,6 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/knowledge-spaces/{id}/source-connections",
"json",
resource_resolver="knowledge_space",
billing_cost=5,
max_request_bytes=256 * 1024,
max_response_bytes=512 * 1024,
stream_kind="json",
@ -614,7 +536,6 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/knowledge-spaces/{id}/source-connections",
"json",
resource_resolver="knowledge_space",
billing_cost=1,
max_request_bytes=16 * 1024,
max_response_bytes=1024 * 1024,
stream_kind="json",
@ -626,7 +547,6 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/knowledge-spaces/{id}/source-connections/{connectionId}/refresh",
"json",
resource_resolver="knowledge_space",
billing_cost=3,
max_request_bytes=32 * 1024,
max_response_bytes=512 * 1024,
stream_kind="json",
@ -638,11 +558,9 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/knowledge-spaces/{id}/sources/{sourceId}/crawl-preview",
"json",
resource_resolver="source",
billing_cost=8,
max_request_bytes=16 * 1024,
max_response_bytes=512 * 1024,
stream_kind="json",
rate_limit_bucket="import",
),
"importSourceWorkflow": _operation(
"POST",
@ -651,11 +569,9 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/knowledge-spaces/{id}/sources/{sourceId}/workflow-imports",
"json",
resource_resolver="source",
billing_cost=25,
max_request_bytes=4 * 1024 * 1024,
max_response_bytes=512 * 1024,
stream_kind="json",
rate_limit_bucket="import",
),
"getSourceSyncPolicy": _operation(
"GET",
@ -664,7 +580,6 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/knowledge-spaces/{id}/sources/{sourceId}/sync-policy",
"json",
resource_resolver="source",
billing_cost=1,
max_request_bytes=0,
max_response_bytes=256 * 1024,
stream_kind="json",
@ -676,7 +591,6 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/knowledge-spaces/{id}/sources/{sourceId}/sync-policy",
"json",
resource_resolver="source",
billing_cost=3,
max_request_bytes=32 * 1024,
max_response_bytes=256 * 1024,
stream_kind="json",
@ -688,7 +602,6 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/knowledge-spaces/{id}/source-workflows/{runId}",
"json",
resource_resolver="job",
billing_cost=1,
max_request_bytes=0,
max_response_bytes=512 * 1024,
stream_kind="json",
@ -700,7 +613,6 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/knowledge-spaces/{id}/source-workflows/{runId}/cancel",
"json",
resource_resolver="job",
billing_cost=2,
max_request_bytes=32 * 1024,
max_response_bytes=512 * 1024,
stream_kind="json",
@ -712,7 +624,6 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/knowledge-spaces/{id}/source-workflows/{runId}/retry",
"json",
resource_resolver="job",
billing_cost=3,
max_request_bytes=0,
max_response_bytes=512 * 1024,
stream_kind="json",
@ -724,7 +635,6 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/knowledge-spaces/{id}/source-workflows/{runId}/pages",
"json",
resource_resolver="job",
billing_cost=1,
max_request_bytes=16 * 1024,
max_response_bytes=4 * 1024 * 1024,
stream_kind="json",
@ -736,11 +646,9 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/knowledge-spaces/{id}/source-workflows/{runId}/selection",
"json",
resource_resolver="job",
billing_cost=8,
max_request_bytes=256 * 1024,
max_response_bytes=512 * 1024,
stream_kind="json",
rate_limit_bucket="import",
),
"crawlSource": _operation(
"POST",
@ -749,11 +657,9 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/knowledge-spaces/{id}/sources/{sourceId}/crawl",
"json",
resource_resolver="source",
billing_cost=20,
max_request_bytes=0,
max_response_bytes=8 * 1024 * 1024,
stream_kind="json",
rate_limit_bucket="import",
),
"listSourcePages": _operation(
"GET",
@ -762,7 +668,6 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/knowledge-spaces/{id}/sources/{sourceId}/pages",
"json",
resource_resolver="source",
billing_cost=3,
max_request_bytes=16 * 1024,
max_response_bytes=4 * 1024 * 1024,
stream_kind="json",
@ -774,11 +679,9 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/knowledge-spaces/{id}/sources/{sourceId}/import",
"json",
resource_resolver="source",
billing_cost=25,
max_request_bytes=1024 * 1024,
max_response_bytes=4 * 1024 * 1024,
stream_kind="json",
rate_limit_bucket="import",
),
"listSourceFiles": _operation(
"GET",
@ -787,7 +690,6 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/knowledge-spaces/{id}/sources/{sourceId}/files",
"json",
resource_resolver="source",
billing_cost=3,
max_request_bytes=32 * 1024,
max_response_bytes=4 * 1024 * 1024,
stream_kind="json",
@ -799,11 +701,9 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/knowledge-spaces/{id}/sources/{sourceId}/import-files",
"json",
resource_resolver="source",
billing_cost=25,
max_request_bytes=1024 * 1024,
max_response_bytes=4 * 1024 * 1024,
stream_kind="json",
rate_limit_bucket="import",
),
"createQuery": _operation(
"POST",
@ -812,7 +712,6 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/queries",
"direct",
resource_resolver="knowledge_space",
billing_cost=20,
max_request_bytes=64 * 1024,
max_response_bytes=0,
stream_kind="sse",
@ -824,7 +723,6 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/knowledge-spaces/{id}/research-tasks",
"json",
resource_resolver="knowledge_space",
billing_cost=2,
max_request_bytes=16 * 1024,
max_response_bytes=2 * 1024 * 1024,
stream_kind="json",
@ -836,7 +734,6 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/research-tasks",
"json",
resource_resolver="knowledge_space",
billing_cost=25,
max_request_bytes=64 * 1024,
max_response_bytes=512 * 1024,
stream_kind="json",
@ -848,7 +745,6 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/research-tasks/plan",
"json",
resource_resolver="knowledge_space",
billing_cost=8,
max_request_bytes=64 * 1024,
max_response_bytes=512 * 1024,
stream_kind="json",
@ -860,7 +756,6 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/research-tasks/{id}",
"json",
resource_resolver="research_task",
billing_cost=1,
max_request_bytes=16 * 1024,
max_response_bytes=512 * 1024,
stream_kind="json",
@ -872,7 +767,6 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/research-tasks/{id}/partials",
"json",
resource_resolver="research_task",
billing_cost=3,
max_request_bytes=16 * 1024,
max_response_bytes=8 * 1024 * 1024,
stream_kind="json",
@ -884,7 +778,6 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/research-tasks/{id}",
"json",
resource_resolver="research_task",
billing_cost=5,
max_request_bytes=16 * 1024,
max_response_bytes=512 * 1024,
stream_kind="json",
@ -896,7 +789,6 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/knowledge-spaces/{id}/golden-questions",
"json",
resource_resolver="knowledge_space",
billing_cost=2,
max_request_bytes=16 * 1024,
max_response_bytes=2 * 1024 * 1024,
stream_kind="json",
@ -908,7 +800,6 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/knowledge-spaces/{id}/golden-questions",
"json",
resource_resolver="knowledge_space",
billing_cost=4,
max_request_bytes=64 * 1024,
max_response_bytes=512 * 1024,
stream_kind="json",
@ -920,7 +811,6 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/knowledge-spaces/{id}/golden-questions/{questionId}",
"json",
resource_resolver="knowledge_space",
billing_cost=4,
max_request_bytes=64 * 1024,
max_response_bytes=512 * 1024,
stream_kind="json",
@ -932,7 +822,6 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/knowledge-spaces/{id}/golden-questions/{questionId}",
"json",
resource_resolver="knowledge_space",
billing_cost=3,
max_request_bytes=0,
max_response_bytes=64 * 1024,
stream_kind="json",
@ -944,7 +833,6 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/knowledge-spaces/{id}/quality/bad-cases",
"json",
resource_resolver="knowledge_space",
billing_cost=2,
max_request_bytes=16 * 1024,
max_response_bytes=2 * 1024 * 1024,
stream_kind="json",
@ -956,7 +844,6 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/knowledge-spaces/{id}/quality/bad-cases",
"json",
resource_resolver="knowledge_space",
billing_cost=4,
max_request_bytes=64 * 1024,
max_response_bytes=512 * 1024,
stream_kind="json",
@ -968,7 +855,6 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/knowledge-spaces/{id}/quality/bad-cases/{badCaseId}",
"json",
resource_resolver="knowledge_space",
billing_cost=2,
max_request_bytes=0,
max_response_bytes=512 * 1024,
stream_kind="json",
@ -980,7 +866,6 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/knowledge-spaces/{id}/quality/bad-cases/{badCaseId}",
"json",
resource_resolver="knowledge_space",
billing_cost=4,
max_request_bytes=64 * 1024,
max_response_bytes=512 * 1024,
stream_kind="json",
@ -992,7 +877,6 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/knowledge-spaces/{id}/quality/replay-runs",
"json",
resource_resolver="knowledge_space",
billing_cost=8,
max_request_bytes=64 * 1024,
max_response_bytes=2 * 1024 * 1024,
stream_kind="json",
@ -1004,7 +888,6 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/knowledge-spaces/{id}/quality/bad-cases/{badCaseId}/trace-reference",
"json",
resource_resolver="knowledge_space",
billing_cost=2,
max_request_bytes=0,
max_response_bytes=64 * 1024,
stream_kind="json",
@ -1016,7 +899,6 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/knowledge-spaces/{id}/quality/traces",
"json",
resource_resolver="knowledge_space",
billing_cost=3,
max_request_bytes=16 * 1024,
max_response_bytes=4 * 1024 * 1024,
stream_kind="json",
@ -1028,7 +910,6 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/queries/{traceId}",
"json",
resource_resolver="query",
billing_cost=2,
max_request_bytes=16 * 1024,
max_response_bytes=2 * 1024 * 1024,
stream_kind="json",
@ -1040,7 +921,6 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/queries/{traceId}/evidence",
"json",
resource_resolver="query",
billing_cost=3,
max_request_bytes=16 * 1024,
max_response_bytes=4 * 1024 * 1024,
stream_kind="json",
@ -1052,7 +932,6 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/queries/{traceId}/conflicts",
"json",
resource_resolver="query",
billing_cost=3,
max_request_bytes=16 * 1024,
max_response_bytes=4 * 1024 * 1024,
stream_kind="json",
@ -1064,7 +943,6 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/queries/{traceId}/missing",
"json",
resource_resolver="query",
billing_cost=3,
max_request_bytes=16 * 1024,
max_response_bytes=4 * 1024 * 1024,
stream_kind="json",
@ -1076,7 +954,6 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/knowledge-spaces/{id}/upload-sessions",
"direct",
resource_resolver="knowledge_space",
billing_cost=5,
max_request_bytes=64 * 1024,
max_response_bytes=64 * 1024,
stream_kind="direct-upload",
@ -1088,7 +965,6 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/upload-sessions/{id}/parts/{partNumber}/presign",
"direct",
resource_resolver="upload_session",
billing_cost=1,
max_request_bytes=32 * 1024,
max_response_bytes=64 * 1024,
stream_kind="direct-upload",
@ -1100,7 +976,6 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/upload-sessions/{id}/small-file",
"binary",
resource_resolver="upload_session",
billing_cost=8,
max_request_bytes=8 * 1024 * 1024,
max_response_bytes=128 * 1024,
stream_kind="json",
@ -1112,7 +987,6 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/upload-sessions/{id}/complete",
"direct",
resource_resolver="upload_session",
billing_cost=8,
max_request_bytes=128 * 1024,
max_response_bytes=128 * 1024,
stream_kind="direct-upload",
@ -1124,7 +998,6 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/upload-sessions/{id}/abort",
"direct",
resource_resolver="upload_session",
billing_cost=1,
max_request_bytes=32 * 1024,
max_response_bytes=64 * 1024,
stream_kind="direct-upload",
@ -1136,7 +1009,6 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO
"/research-tasks/{id}/events",
"direct",
resource_resolver="research_task",
billing_cost=5,
max_request_bytes=0,
max_response_bytes=0,
stream_kind="sse",

View File

@ -28,13 +28,6 @@ from services.knowledge_fs.cutover import KnowledgeFSWorkspaceCutoverService
from services.knowledge_fs.cutover_runtime_gate import SQLKnowledgeFSWorkspaceRuntimeGate
from services.knowledge_fs.data_facade import KnowledgeFSDataFacade
from services.knowledge_fs.greenfield_initializer import KnowledgeFSWorkspaceGreenfieldInitializer
from services.knowledge_fs.operation_admission import (
DifyKnowledgeFSBillingPort,
DifyKnowledgeFSWeightedRateLimitPort,
KnowledgeFSDirectOperationAdmissionService,
KnowledgeFSOperationAdmissionService,
LoggingKnowledgeFSRateLimitAudit,
)
from services.knowledge_fs.product_application_service import KnowledgeFSProductApplicationService
from services.knowledge_fs.product_authorization import DifyKnowledgeFSProductRBACPort
from services.knowledge_fs.product_remote import KnowledgeFSOperationUnavailableError
@ -53,9 +46,7 @@ class KnowledgeFSRuntime(NamedTuple):
broker: KnowledgeFSCapabilityBroker
control_plane: KnowledgeFSControlPlaneService
credentials: KnowledgeFSCredentialService
direct_operation_admission: KnowledgeFSDirectOperationAdmissionService
facade: KnowledgeFSDataFacade
operation_admission: KnowledgeFSOperationAdmissionService
def create_knowledge_fs_runtime(session_maker: sessionmaker[Session]) -> KnowledgeFSRuntime:
@ -108,15 +99,7 @@ def create_knowledge_fs_runtime(session_maker: sessionmaker[Session]) -> Knowled
product=product,
issuer=issuer,
)
operation_admission = KnowledgeFSOperationAdmissionService(
rate_limit=DifyKnowledgeFSWeightedRateLimitPort(audit=LoggingKnowledgeFSRateLimitAudit()),
billing=DifyKnowledgeFSBillingPort(),
)
direct_operation_admission = KnowledgeFSDirectOperationAdmissionService(
admission=operation_admission,
broker=broker,
)
facade = KnowledgeFSDataFacade(admission=operation_admission, broker=broker, remote=remote)
facade = KnowledgeFSDataFacade(broker=broker, remote=remote)
credentials = KnowledgeFSCredentialService(session_maker, product=product, revocations=revocations)
application = KnowledgeFSProductApplicationService(
product=product,
@ -138,15 +121,12 @@ def create_knowledge_fs_runtime(session_maker: sessionmaker[Session]) -> Knowled
app_capabilities=KnowledgeFSAppExecutionCapabilityService(
admission=app_admission,
broker=broker,
operation_admission=operation_admission,
remote=remote,
),
broker=broker,
control_plane=control_plane,
credentials=credentials,
direct_operation_admission=direct_operation_admission,
facade=facade,
operation_admission=operation_admission,
)

View File

@ -6,25 +6,18 @@ from datetime import UTC, datetime
from io import BytesIO
from pathlib import Path
from types import SimpleNamespace
from unittest.mock import MagicMock
import pytest
from flask import Flask
from werkzeug.exceptions import RequestEntityTooLarge
from werkzeug.exceptions import Forbidden, RequestEntityTooLarge
from controllers.console import console_ns
from controllers.console import wraps as console_wraps
from controllers.console.knowledge_fs import resources as console_resources
from controllers.console.knowledge_fs.error import KnowledgeFSQuotaExceededHTTPError, KnowledgeFSRateLimitHTTPError
from controllers.service_api import service_api_ns
from controllers.service_api.knowledge_fs import resources as service_resources
from controllers.service_api.knowledge_fs.error import (
KnowledgeFSServiceQuotaExceededHTTPError,
KnowledgeFSServiceRateLimitHTTPError,
)
from services.knowledge_fs.credential_service import KnowledgeFSServiceCredentialProfile
from services.knowledge_fs.operation_admission import (
KnowledgeFSOperationQuotaExceededError,
KnowledgeFSOperationRateLimitExceededError,
)
from services.knowledge_fs.product_dto import (
KnowledgeFSDocumentUploadAcceptedResponse,
KnowledgeFSSmallFileUploadResponse,
@ -312,7 +305,13 @@ def test_document_upload_console_bff_reads_only_through_facade_and_returns_accep
]
def test_small_file_console_bff_maps_oversize_to_413() -> None:
def test_small_file_console_bff_maps_oversize_to_413(monkeypatch: pytest.MonkeyPatch) -> None:
monkeypatch.setattr(console_wraps, "current_account_with_tenant", lambda: (object(), "tenant-1"))
monkeypatch.setattr(
console_wraps.FeatureService,
"get_knowledge_rate_limit",
lambda _tenant_id: SimpleNamespace(enabled=False),
)
app = Flask(__name__)
with app.test_request_context(
method="POST",
@ -331,35 +330,31 @@ def test_small_file_console_bff_maps_oversize_to_413() -> None:
reject()
def test_console_maps_weighted_rate_and_billing_exhaustion_to_stable_errors() -> None:
@console_resources._knowledge_fs_errors
def rate_limited():
raise KnowledgeFSOperationRateLimitExceededError()
def test_console_knowledge_fs_uses_existing_knowledge_rate_limit(monkeypatch: pytest.MonkeyPatch) -> None:
redis = MagicMock()
redis.zcard.return_value = 2
session = MagicMock()
monkeypatch.setattr(console_wraps, "current_account_with_tenant", lambda: (object(), "tenant-1"))
monkeypatch.setattr(
console_wraps.FeatureService,
"get_knowledge_rate_limit",
lambda _tenant_id: SimpleNamespace(enabled=True, limit=1, subscription_plan="sandbox"),
)
monkeypatch.setattr(console_wraps, "redis_client", redis)
monkeypatch.setattr(console_wraps, "db", SimpleNamespace(session=session))
@console_resources._knowledge_fs_errors
def quota_exhausted():
raise KnowledgeFSOperationQuotaExceededError()
def view():
return "should not run"
with pytest.raises(KnowledgeFSRateLimitHTTPError):
rate_limited()
with pytest.raises(KnowledgeFSQuotaExceededHTTPError):
quota_exhausted()
app = Flask(__name__)
with app.test_request_context(), pytest.raises(Forbidden, match="knowledge base request rate limit"):
view()
@service_resources._service_api_errors
def service_rate_limited():
raise KnowledgeFSOperationRateLimitExceededError()
@service_resources._service_api_errors
def service_quota_exhausted():
raise KnowledgeFSOperationQuotaExceededError()
with pytest.raises(KnowledgeFSServiceRateLimitHTTPError) as service_rate:
service_rate_limited()
with pytest.raises(KnowledgeFSServiceQuotaExceededHTTPError) as service_quota:
service_quota_exhausted()
assert service_rate.value.code == 429
assert service_quota.value.code == 403
redis.zadd.assert_called_once()
redis.zremrangebyscore.assert_called_once()
session.add.assert_called_once()
session.commit.assert_called_once()
def test_space_update_and_delete_publish_their_actual_http_status_contracts() -> None:
@ -571,7 +566,7 @@ def test_query_stream_capability_issues_exact_space_grant_without_token_in_url(
) -> None:
calls: list[dict[str, object]] = []
class DirectAdmission:
class Broker:
def issue_interactive(self, **kwargs):
calls.append(kwargs)
return SimpleNamespace(
@ -584,7 +579,7 @@ def test_query_stream_capability_issues_exact_space_grant_without_token_in_url(
monkeypatch.setattr(
console_resources,
"_console_services",
lambda: SimpleNamespace(direct_operation_admission=DirectAdmission()),
lambda: SimpleNamespace(broker=Broker()),
)
app = Flask(__name__)
@ -613,7 +608,7 @@ def test_query_stream_capability_issues_exact_space_grant_without_token_in_url(
def test_query_admission_binds_validated_mode_to_resolved_kfs_space(monkeypatch: pytest.MonkeyPatch) -> None:
calls: list[dict[str, object]] = []
class DirectAdmission:
class Broker:
def issue_interactive(self, **kwargs):
calls.append(kwargs)
return SimpleNamespace(
@ -627,7 +622,7 @@ def test_query_admission_binds_validated_mode_to_resolved_kfs_space(monkeypatch:
monkeypatch.setattr(
console_resources,
"_console_services",
lambda: SimpleNamespace(direct_operation_admission=DirectAdmission()),
lambda: SimpleNamespace(broker=Broker()),
)
app = Flask(__name__)
@ -654,12 +649,12 @@ def test_query_admission_binds_validated_mode_to_resolved_kfs_space(monkeypatch:
]
def test_upload_and_task_stream_capabilities_use_direct_operation_admission(
def test_upload_and_task_stream_capabilities_use_broker(
monkeypatch: pytest.MonkeyPatch,
) -> None:
calls: list[dict[str, object]] = []
class DirectAdmission:
class Broker:
def issue_interactive(self, **kwargs):
calls.append(kwargs)
return SimpleNamespace(
@ -668,7 +663,7 @@ def test_upload_and_task_stream_capabilities_use_direct_operation_admission(
knowledge_space_id="space-1",
)
runtime = SimpleNamespace(direct_operation_admission=DirectAdmission())
runtime = SimpleNamespace(broker=Broker())
monkeypatch.setattr(console_resources.dify_config, "KNOWLEDGE_FS_DIRECT_ORIGIN", "https://kfs.test")
monkeypatch.setattr(console_resources.dify_config, "KNOWLEDGE_FS_DIRECT_UPLOAD_READY", True)
monkeypatch.setattr(console_resources, "_actor", lambda: ("account-1", "tenant-1"))
@ -703,7 +698,7 @@ def test_upload_and_task_stream_capabilities_use_direct_operation_admission(
]
def test_service_query_admission_uses_direct_operation_admission(monkeypatch: pytest.MonkeyPatch) -> None:
def test_service_query_admission_uses_broker(monkeypatch: pytest.MonkeyPatch) -> None:
calls: list[dict[str, object]] = []
profile = SimpleNamespace(tenant_id="tenant-1", control_space_id="control-1")
@ -712,7 +707,7 @@ def test_service_query_admission_uses_direct_operation_admission(monkeypatch: py
_ = kwargs
return profile
class DirectAdmission:
class Broker:
def issue_service(self, **kwargs):
calls.append(kwargs)
return SimpleNamespace(
@ -723,7 +718,7 @@ def test_service_query_admission_uses_direct_operation_admission(monkeypatch: py
runtime = SimpleNamespace(
credentials=Credentials(),
direct_operation_admission=DirectAdmission(),
broker=Broker(),
)
monkeypatch.setattr(service_resources.dify_config, "KNOWLEDGE_FS_DIRECT_ORIGIN", "https://kfs.test")
monkeypatch.setattr(service_resources, "_runtime", lambda: runtime)

View File

@ -16,6 +16,7 @@ from werkzeug.exceptions import (
UnprocessableEntity,
)
from controllers.console import wraps as console_wraps
from controllers.console.knowledge_fs import resources as console_resources
from controllers.console.knowledge_fs.error import (
KnowledgeFSAccessDeniedHTTPError,
@ -651,7 +652,7 @@ def test_console_overview_stats_composes_kfs_metrics_with_dify_app_bindings(
),
app_bindings=SimpleNamespace(count_active=MagicMock(return_value=7)),
)
dump_response = MagicMock(side_effect=lambda schema, raw: raw)
dump_response = MagicMock(side_effect=lambda _schema, raw: raw)
monkeypatch.setattr(console_resources, "_actor", lambda: ("account-1", "tenant-1"))
monkeypatch.setattr(console_resources, "_console_services", lambda: runtime)
monkeypatch.setattr(console_resources, "dump_response", dump_response)
@ -985,7 +986,7 @@ def test_console_direct_capabilities_bind_the_authorized_resource(monkeypatch: p
expires_at=datetime(2026, 7, 21, tzinfo=UTC),
knowledge_space_id="knowledge-space-1",
)
admission = SimpleNamespace(issue_interactive=MagicMock(return_value=issued))
broker = SimpleNamespace(issue_interactive=MagicMock(return_value=issued))
payloads = iter(
[
KnowledgeFSUploadCapabilityPayload(operation_id="completeUploadSession", upload_session_id="session-1"),
@ -999,7 +1000,7 @@ def test_console_direct_capabilities_bind_the_authorized_resource(monkeypatch: p
monkeypatch.setattr(
console_resources,
"_console_services",
lambda: SimpleNamespace(direct_operation_admission=admission),
lambda: SimpleNamespace(broker=broker),
)
monkeypatch.setattr(console_resources, "_payload", lambda _: next(payloads))
monkeypatch.setattr(console_resources, "dump_response", lambda _, response: response)
@ -1016,9 +1017,9 @@ def test_console_direct_capabilities_bind_the_authorized_resource(monkeypatch: p
assert query.url == "https://kfs.example/queries"
assert task.url == "https://kfs.example//research-tasks/task%2F1/events?knowledgeSpaceId=knowledge-space-1"
assert legacy_query.url == "https://kfs.example/queries"
assert admission.issue_interactive.call_args_list[0].kwargs["resource_id"] == "session-1"
assert admission.issue_interactive.call_args_list[1].kwargs["operation_id"] == "createQuery"
assert admission.issue_interactive.call_args_list[2].kwargs == {
assert broker.issue_interactive.call_args_list[0].kwargs["resource_id"] == "session-1"
assert broker.issue_interactive.call_args_list[1].kwargs["operation_id"] == "createQuery"
assert broker.issue_interactive.call_args_list[2].kwargs == {
"tenant_id": "tenant-1",
"account_id": "account-1",
"control_space_id": "space-2",
@ -1033,8 +1034,8 @@ def test_service_direct_query_admission_binds_profile_space_and_payload(monkeypa
expires_at=datetime(2026, 7, 21, tzinfo=UTC),
knowledge_space_id="knowledge-space-1",
)
admission = SimpleNamespace(issue_service=MagicMock(return_value=issued))
runtime = SimpleNamespace(direct_operation_admission=admission)
broker = SimpleNamespace(issue_service=MagicMock(return_value=issued))
runtime = SimpleNamespace(broker=broker)
profile = object()
monkeypatch.setattr(service_resources.dify_config, "KNOWLEDGE_FS_DIRECT_ORIGIN", "https://kfs.example/")
monkeypatch.setattr(service_resources, "_runtime", lambda: runtime)
@ -1046,7 +1047,7 @@ def test_service_direct_query_admission_binds_profile_space_and_payload(monkeypa
with app.test_request_context("/", method="POST"):
response = _invoke(service_resources, "KnowledgeFSServiceQueryAdmissionApi", "post", "space-1")
admission.issue_service.assert_called_once_with(profile=profile, operation_id="createQuery")
broker.issue_service.assert_called_once_with(profile=profile, operation_id="createQuery")
assert response.request.knowledge_space_id == "knowledge-space-1"
assert response.request.query == "question"
assert response.url == "https://kfs.example/queries"
@ -1164,9 +1165,17 @@ def test_service_resource_helpers_validate_feature_bearer_headers_and_boolean_qu
service_resources._runtime()
def test_console_request_rejections_preserve_conflict_size_and_validation_contracts() -> None:
def test_console_request_rejections_preserve_conflict_size_and_validation_contracts(
monkeypatch: pytest.MonkeyPatch,
) -> None:
from services.knowledge_fs.product_remote import KnowledgeFSProductRequestRejectedError
monkeypatch.setattr(console_wraps, "current_account_with_tenant", lambda: (object(), "tenant-1"))
monkeypatch.setattr(
console_wraps.FeatureService,
"get_knowledge_rate_limit",
lambda _tenant_id: SimpleNamespace(enabled=False),
)
expected = {
HTTPStatus.CONFLICT: Conflict,
HTTPStatus.REQUEST_ENTITY_TOO_LARGE: RequestEntityTooLarge,
@ -1205,7 +1214,9 @@ def test_jwks_resource_fails_closed_for_disabled_missing_and_misconfigured_issue
console_resources.KnowledgeFSJWKSApi().get()
def test_console_error_adapter_maps_every_domain_boundary_to_the_stable_http_contract() -> None:
def test_console_error_adapter_maps_every_domain_boundary_to_the_stable_http_contract(
monkeypatch: pytest.MonkeyPatch,
) -> None:
from pydantic import ValidationError
from services.knowledge_fs.app_binding_management import KnowledgeFSAppBindingManagementError
@ -1218,6 +1229,12 @@ def test_console_error_adapter_maps_every_domain_boundary_to_the_stable_http_con
)
from services.knowledge_fs_capability import KnowledgeFSCapabilityConfigurationError
monkeypatch.setattr(console_wraps, "current_account_with_tenant", lambda: (object(), "tenant-1"))
monkeypatch.setattr(
console_wraps.FeatureService,
"get_knowledge_rate_limit",
lambda _tenant_id: SimpleNamespace(enabled=False),
)
with pytest.raises(ValidationError) as raised_validation:
KnowledgeFSQueryCreatePayload.model_validate({"query": ""})
validation_error = raised_validation.value

View File

@ -1,7 +1,6 @@
from __future__ import annotations
from datetime import UTC, datetime
from typing import Any
import pytest
from sqlalchemy.orm import Session, sessionmaker
@ -65,7 +64,7 @@ class Remote:
def __init__(self) -> None:
self.calls: list[KnowledgeFSRemoteJSONRequest] = []
def execute_json(self, request: KnowledgeFSRemoteJSONRequest) -> dict[str, Any]:
def execute_json(self, request: KnowledgeFSRemoteJSONRequest) -> dict[str, object]:
self.calls.append(request)
return {
"id": "task-1",
@ -79,28 +78,6 @@ class Remote:
}
class Charge:
def __init__(self) -> None:
self.committed = False
self.refunded = False
def commit(self) -> None:
self.committed = True
def refund(self) -> None:
self.refunded = True
class OperationAdmission:
def __init__(self) -> None:
self.charge = Charge()
def reserve(self, *, tenant_id: str, operation_id: str) -> Charge:
assert tenant_id == "tenant-1"
assert operation_id == "createResearchTask"
return self.charge
def _run_context(*, app_id: str) -> DifyRunContext:
return DifyRunContext(
tenant_id="tenant-1",
@ -198,11 +175,9 @@ def test_execution_fails_closed_before_capability_or_remote_io(
)
broker = Broker()
remote = Remote()
operation_admission = OperationAdmission()
service = KnowledgeFSAppExecutionCapabilityService( # type: ignore[arg-type]
admission=KnowledgeFSAppAdmissionService(sessionmaker(bind=sqlite_session.get_bind(), expire_on_commit=False)),
broker=broker,
operation_admission=operation_admission,
remote=remote,
)
@ -216,8 +191,6 @@ def test_execution_fails_closed_before_capability_or_remote_io(
assert broker.calls == []
assert remote.calls == []
assert operation_admission.charge.committed is False
assert operation_admission.charge.refunded is True
@pytest.mark.parametrize(
@ -231,11 +204,9 @@ def test_execution_issues_capability_then_performs_bounded_remote_io(
space = _seed(sqlite_session, caller_kind=caller_kind)
broker = Broker()
remote = Remote()
operation_admission = OperationAdmission()
service = KnowledgeFSAppExecutionCapabilityService( # type: ignore[arg-type]
admission=KnowledgeFSAppAdmissionService(sessionmaker(bind=sqlite_session.get_bind(), expire_on_commit=False)),
broker=broker,
operation_admission=operation_admission,
remote=remote,
)
@ -248,8 +219,6 @@ def test_execution_issues_capability_then_performs_bounded_remote_io(
assert response.id == "task-1"
assert len(broker.calls) == 1
assert operation_admission.charge.committed is True
assert operation_admission.charge.refunded is False
assert remote.calls == [
KnowledgeFSRemoteJSONRequest(
operation_id="createResearchTask",

View File

@ -2,20 +2,20 @@ from __future__ import annotations
from datetime import UTC, datetime
from types import SimpleNamespace
from typing import Any
import pytest
from pydantic import ValidationError
from core.app.entities.app_invoke_entities import DifyRunContext, InvokeFrom, UserFrom
from models.knowledge_fs import KnowledgeFSAppSpaceJoinType
from services.knowledge_fs import app_execution_capability
from services.knowledge_fs.app_execution_capability import (
KnowledgeFSAppExecutionCapabilityService,
KnowledgeResourceRef,
)
from services.knowledge_fs.capability_broker import KnowledgeFSIssuedProductCapability
from services.knowledge_fs.product_dto import KnowledgeFSResearchTaskCreatePayload
from services.knowledge_fs.product_remote import KnowledgeFSRemoteJSONRequest
from services.knowledge_fs.product_remote import KnowledgeFSOperationUnavailableError, KnowledgeFSRemoteJSONRequest
class Admission:
@ -48,7 +48,7 @@ class Remote:
def __init__(self) -> None:
self.calls: list[KnowledgeFSRemoteJSONRequest] = []
def execute_json(self, request: KnowledgeFSRemoteJSONRequest) -> dict[str, Any]:
def execute_json(self, request: KnowledgeFSRemoteJSONRequest) -> dict[str, object]:
self.calls.append(request)
return {
"id": "research-1",
@ -62,47 +62,18 @@ class Remote:
}
class Charge:
def __init__(self) -> None:
self.committed = False
self.refunded = False
def commit(self) -> None:
self.committed = True
def refund(self) -> None:
self.refunded = True
class OperationAdmission:
def __init__(self) -> None:
self.calls: list[dict[str, str]] = []
self.charge = Charge()
def reserve(self, *, tenant_id: str, operation_id: str) -> Charge:
self.calls.append({"tenant_id": tenant_id, "operation_id": operation_id})
return self.charge
class FailingRemote(Remote):
def execute_json(self, request: KnowledgeFSRemoteJSONRequest) -> dict[str, Any]:
def execute_json(self, request: KnowledgeFSRemoteJSONRequest) -> dict[str, object]:
self.calls.append(request)
raise RuntimeError("remote failed")
class RejectingOperationAdmission(OperationAdmission):
def reserve(self, *, tenant_id: str, operation_id: str) -> Charge:
self.calls.append({"tenant_id": tenant_id, "operation_id": operation_id})
raise RuntimeError("operation rejected")
def test_app_execution_capability_always_admits_binding_before_broker_issuance() -> None:
admission = Admission()
broker = Broker()
service = KnowledgeFSAppExecutionCapabilityService( # type: ignore[arg-type]
admission=admission,
broker=broker,
operation_admission=OperationAdmission(),
remote=Remote(),
)
@ -154,16 +125,54 @@ def test_typed_knowledge_resource_ref_rejects_dataset_and_extra_authority_fields
}
)
with pytest.raises(ValidationError):
KnowledgeResourceRef.model_validate(
{
"kind": "knowledge_fs",
"control_space_id": " ",
}
)
def test_create_research_task_fails_closed_when_product_operation_is_unavailable(
monkeypatch: pytest.MonkeyPatch,
) -> None:
admission = Admission()
broker = Broker()
remote = Remote()
service = KnowledgeFSAppExecutionCapabilityService( # type: ignore[arg-type]
admission=admission,
broker=broker,
remote=remote,
)
monkeypatch.setattr(app_execution_capability, "is_product_operation_ready", lambda _operation_id: False)
with pytest.raises(KnowledgeFSOperationUnavailableError, match="creation is unavailable"):
service.create_research_task(
run_context=DifyRunContext(
tenant_id="tenant-1",
app_id="app-1",
user_id="user-1",
user_from=UserFrom.END_USER,
invoke_from=InvokeFrom.WEB_APP,
),
caller_kind=KnowledgeFSAppSpaceJoinType.AGENT,
resource=KnowledgeResourceRef(kind="knowledge_fs", control_space_id="control-1"),
payload=KnowledgeFSResearchTaskCreatePayload(query="Compare the evidence"),
)
assert admission.calls == []
assert broker.calls == []
assert remote.calls == []
def test_create_research_task_uses_only_dify_run_context_identity_then_calls_kfs() -> None:
admission = Admission()
broker = Broker()
remote = Remote()
operation_admission = OperationAdmission()
service = KnowledgeFSAppExecutionCapabilityService( # type: ignore[arg-type]
admission=admission,
broker=broker,
operation_admission=operation_admission,
remote=remote,
)
run_context = DifyRunContext(
@ -216,22 +225,12 @@ def test_create_research_task_uses_only_dify_run_context_identity_then_calls_kfs
},
)
]
assert operation_admission.calls == [
{
"tenant_id": "tenant-from-run-context",
"operation_id": "createResearchTask",
}
]
assert operation_admission.charge.committed is True
assert operation_admission.charge.refunded is False
def test_create_research_task_refunds_operation_charge_when_remote_io_fails() -> None:
operation_admission = OperationAdmission()
def test_create_research_task_propagates_remote_io_failure() -> None:
service = KnowledgeFSAppExecutionCapabilityService( # type: ignore[arg-type]
admission=Admission(),
broker=Broker(),
operation_admission=operation_admission,
remote=FailingRemote(),
)
@ -248,37 +247,3 @@ def test_create_research_task_refunds_operation_charge_when_remote_io_fails() ->
resource=KnowledgeResourceRef(kind="knowledge_fs", control_space_id="control-1"),
payload=KnowledgeFSResearchTaskCreatePayload(query="Compare the evidence"),
)
assert operation_admission.charge.committed is False
assert operation_admission.charge.refunded is True
def test_operation_admission_rejection_prevents_app_capability_and_remote_io() -> None:
admission = Admission()
broker = Broker()
remote = Remote()
operation_admission = RejectingOperationAdmission()
service = KnowledgeFSAppExecutionCapabilityService( # type: ignore[arg-type]
admission=admission,
broker=broker,
operation_admission=operation_admission,
remote=remote,
)
with pytest.raises(RuntimeError, match="operation rejected"):
service.create_research_task(
run_context=DifyRunContext(
tenant_id="tenant-1",
app_id="app-1",
user_id="user-1",
user_from=UserFrom.END_USER,
invoke_from=InvokeFrom.WEB_APP,
),
caller_kind=KnowledgeFSAppSpaceJoinType.AGENT,
resource=KnowledgeResourceRef(kind="knowledge_fs", control_space_id="control-1"),
payload=KnowledgeFSResearchTaskCreatePayload(query="Compare the evidence"),
)
assert admission.calls == []
assert broker.calls == []
assert remote.calls == []

View File

@ -31,44 +31,6 @@ from services.knowledge_fs.product_remote import (
)
class NoopCharge:
def commit(self) -> None:
return
def refund(self) -> None:
return
class NoopAdmission:
def reserve(self, **kwargs: object) -> NoopCharge:
_ = kwargs
return NoopCharge()
class RecordingCharge:
def __init__(self) -> None:
self.commits = 0
self.refunds = 0
def commit(self) -> None:
self.commits += 1
def refund(self) -> None:
self.refunds += 1
class RecordingAdmission:
def __init__(self) -> None:
self.calls: list[dict[str, object]] = []
self.charges: list[RecordingCharge] = []
def reserve(self, **kwargs: object) -> RecordingCharge:
self.calls.append(kwargs)
charge = RecordingCharge()
self.charges.append(charge)
return charge
class FailingBroker:
def __init__(self) -> None:
self.calls = 0
@ -327,10 +289,8 @@ class ActiveSettingsRemote(RecordingRemote):
return super().execute_json(request)
def test_facade_reserves_before_capability_and_commits_or_refunds_after_remote_io() -> None:
admission = RecordingAdmission()
def test_facade_propagates_remote_io_failure() -> None:
facade = KnowledgeFSDataFacade( # type: ignore[arg-type]
admission=admission,
broker=RecordingBroker(),
remote=RecordingRemote(),
)
@ -342,12 +302,8 @@ def test_facade_reserves_before_capability_and_commits_or_refunds_after_remote_i
)
assert result.revision == 2
assert admission.calls == [{"tenant_id": "tenant-1", "operation_id": "getSettings"}]
assert (admission.charges[0].commits, admission.charges[0].refunds) == (1, 0)
failing_admission = RecordingAdmission()
failing = KnowledgeFSDataFacade( # type: ignore[arg-type]
admission=failing_admission,
broker=RecordingBroker(),
remote=FailingRemote(),
)
@ -357,13 +313,12 @@ def test_facade_reserves_before_capability_and_commits_or_refunds_after_remote_i
account_id="account-1",
control_space_id="control-1",
)
assert (failing_admission.charges[0].commits, failing_admission.charges[0].refunds) == (0, 1)
def test_document_upload_authorizes_before_read_and_binds_bounded_multipart_request() -> None:
remote = RecordingRemote()
broker = RecordingBroker()
facade = KnowledgeFSDataFacade(admission=NoopAdmission(), broker=broker, remote=remote) # type: ignore[arg-type]
facade = KnowledgeFSDataFacade(broker=broker, remote=remote) # type: ignore[arg-type]
observed: list[str] = []
def read_upload(max_bytes: int) -> KnowledgeFSRemoteMultipartFile:
@ -413,7 +368,7 @@ def test_document_upload_authorizes_before_read_and_binds_bounded_multipart_requ
def test_legacy_buffered_query_fails_before_capability_or_remote_io() -> None:
broker = FailingBroker()
remote = FailingRemote()
facade = KnowledgeFSDataFacade(admission=NoopAdmission(), broker=broker, remote=remote) # type: ignore[arg-type]
facade = KnowledgeFSDataFacade(broker=broker, remote=remote) # type: ignore[arg-type]
with pytest.raises(KnowledgeFSOperationUnavailableError, match="queries/admission"):
facade.create_query(
@ -430,7 +385,7 @@ def test_legacy_buffered_query_fails_before_capability_or_remote_io() -> None:
def test_small_file_fallback_authorizes_before_read_and_binds_narrow_binary_request() -> None:
remote = RecordingRemote()
broker = RecordingBroker()
facade = KnowledgeFSDataFacade(admission=NoopAdmission(), broker=broker, remote=remote) # type: ignore[arg-type]
facade = KnowledgeFSDataFacade(broker=broker, remote=remote) # type: ignore[arg-type]
observed: list[str] = []
def read_body(max_bytes: int) -> bytes:
@ -482,7 +437,8 @@ def test_small_file_fallback_denial_and_size_limit_stop_before_bytes_or_remote_i
denied_remote = FailingRemote()
denied_facade = KnowledgeFSDataFacade( # type: ignore[arg-type]
admission=NoopAdmission(), broker=DenyingBroker(), remote=denied_remote
broker=DenyingBroker(),
remote=denied_remote,
)
reads = 0
@ -504,7 +460,8 @@ def test_small_file_fallback_denial_and_size_limit_stop_before_bytes_or_remote_i
remote = RecordingRemote()
facade = KnowledgeFSDataFacade( # type: ignore[arg-type]
admission=NoopAdmission(), broker=RecordingBroker(), remote=remote
broker=RecordingBroker(),
remote=remote,
)
with pytest.raises(KnowledgeFSProductRequestRejectedError) as oversized:
facade.upload_small_file(
@ -521,7 +478,7 @@ def test_small_file_fallback_denial_and_size_limit_stop_before_bytes_or_remote_i
def test_json_facade_uses_kfs_camel_case_body_and_authoritative_revision() -> None:
remote = RecordingRemote()
broker = RecordingBroker()
facade = KnowledgeFSDataFacade(admission=NoopAdmission(), broker=broker, remote=remote) # type: ignore[arg-type]
facade = KnowledgeFSDataFacade(broker=broker, remote=remote) # type: ignore[arg-type]
research = facade.create_research_task(
tenant_id="tenant-1",
@ -560,7 +517,7 @@ def test_json_facade_uses_kfs_camel_case_body_and_authoritative_revision() -> No
def test_basic_product_facade_resolves_control_space_then_uses_exact_kfs_routes() -> None:
remote = RecordingRemote()
broker = RecordingBroker()
facade = KnowledgeFSDataFacade(admission=NoopAdmission(), broker=broker, remote=remote) # type: ignore[arg-type]
facade = KnowledgeFSDataFacade(broker=broker, remote=remote) # type: ignore[arg-type]
settings = facade.get_settings(
tenant_id="tenant-1",
@ -656,7 +613,6 @@ def test_basic_product_facade_resolves_control_space_then_uses_exact_kfs_routes(
def test_active_settings_use_durable_profile_migration_routes() -> None:
embedding_remote = ActiveSettingsRemote()
embedding_facade = KnowledgeFSDataFacade(
admission=NoopAdmission(),
broker=RecordingBroker(),
remote=embedding_remote,
) # type: ignore[arg-type]
@ -691,7 +647,6 @@ def test_active_settings_use_durable_profile_migration_routes() -> None:
retrieval_remote = ActiveSettingsRemote()
retrieval_facade = KnowledgeFSDataFacade(
admission=NoopAdmission(),
broker=RecordingBroker(),
remote=retrieval_remote,
) # type: ignore[arg-type]
@ -749,7 +704,6 @@ def test_active_settings_use_durable_profile_migration_routes() -> None:
def test_get_profile_migration_uses_the_durable_migration_route() -> None:
remote = ActiveSettingsRemote()
facade = KnowledgeFSDataFacade(
admission=NoopAdmission(),
broker=RecordingBroker(),
remote=remote,
) # type: ignore[arg-type]
@ -768,7 +722,6 @@ def test_get_profile_migration_uses_the_durable_migration_route() -> None:
def test_active_settings_reject_concurrent_profile_migrations() -> None:
remote = ActiveSettingsRemote()
facade = KnowledgeFSDataFacade(
admission=NoopAdmission(),
broker=RecordingBroker(),
remote=remote,
) # type: ignore[arg-type]
@ -829,7 +782,7 @@ def test_settings_dto_rejects_fast_threshold_without_rerank() -> None:
def test_advanced_facade_binds_child_resources_parent_space_and_idempotency() -> None:
remote = RecordingRemote()
broker = RecordingBroker()
facade = KnowledgeFSDataFacade(admission=NoopAdmission(), broker=broker, remote=remote) # type: ignore[arg-type]
facade = KnowledgeFSDataFacade(broker=broker, remote=remote) # type: ignore[arg-type]
source = facade.update_source(
tenant_id="tenant-1",
@ -1230,7 +1183,7 @@ def test_facade_public_methods_preserve_the_registered_operation_and_child_bindi
specific_kwargs: dict[str, object],
child_resource_id: str | None,
) -> None:
facade = KnowledgeFSDataFacade(admission=MagicMock(), broker=MagicMock(), remote=MagicMock())
facade = KnowledgeFSDataFacade(broker=MagicMock(), remote=MagicMock())
interactive = MagicMock(return_value={})
interactive_child = MagicMock(return_value={})
response_type = getattr(data_facade_module, response_name)

View File

@ -7,7 +7,6 @@ from services.knowledge_fs.observability import (
KnowledgeFSCapabilityIssuanceMetric,
KnowledgeFSControlSpaceStateMetric,
KnowledgeFSLifecycleTaskMetric,
KnowledgeFSOperationAdmissionMetric,
OpenTelemetryKnowledgeFSOperationalMetrics,
)
@ -51,15 +50,6 @@ def test_operational_metrics_use_only_bounded_labels_and_numeric_measurements()
status="succeeded",
)
)
metrics.record_operation_admission(
KnowledgeFSOperationAdmissionMetric(
bucket="direct",
operation_id="createQuery",
outcome="success",
phase="commit",
)
)
capability_counter = counters["dify.knowledge_fs.capability_issuance"]
assert capability_counter.add.call_args.args == (1,)
assert capability_counter.add.call_args.kwargs == {

View File

@ -1,310 +0,0 @@
from __future__ import annotations
from types import SimpleNamespace
from unittest.mock import MagicMock
import pytest
from services.knowledge_fs.operation_admission import (
DifyKnowledgeFSBillingPort,
DifyKnowledgeFSWeightedRateLimitPort,
KnowledgeFSDirectOperationAdmissionService,
KnowledgeFSOperationAdmissionService,
KnowledgeFSOperationQuotaExceededError,
KnowledgeFSOperationRateLimitExceededError,
KnowledgeFSOperationUsage,
)
class RecordingCharge:
def __init__(self) -> None:
self.commits = 0
self.refunds = 0
def commit(self) -> None:
self.commits += 1
def refund(self) -> None:
self.refunds += 1
class RecordingRateLimit:
def __init__(self) -> None:
self.usages: list[KnowledgeFSOperationUsage] = []
def admit(self, usage: KnowledgeFSOperationUsage) -> None:
self.usages.append(usage)
class RecordingBilling:
def __init__(self, charge: RecordingCharge) -> None:
self.charge = charge
self.usages: list[KnowledgeFSOperationUsage] = []
def reserve(self, usage: KnowledgeFSOperationUsage) -> RecordingCharge:
self.usages.append(usage)
return self.charge
class RecordingBroker:
def __init__(self, *, error: Exception | None = None) -> None:
self.error = error
self.interactive_calls: list[dict[str, object]] = []
self.service_calls: list[dict[str, object]] = []
def issue_interactive(self, **kwargs):
self.interactive_calls.append(kwargs)
if self.error is not None:
raise self.error
return SimpleNamespace(token="interactive-capability")
def issue_service(self, **kwargs):
self.service_calls.append(kwargs)
if self.error is not None:
raise self.error
return SimpleNamespace(token="service-capability")
def test_operation_admission_consumes_registry_billing_and_rate_limit_mapping() -> None:
rate_limit = RecordingRateLimit()
charge = RecordingCharge()
billing = RecordingBilling(charge)
service = KnowledgeFSOperationAdmissionService(rate_limit=rate_limit, billing=billing)
returned = service.reserve(tenant_id="tenant-1", operation_id="reindexDocuments")
assert returned is charge
assert rate_limit.usages == [
KnowledgeFSOperationUsage(
tenant_id="tenant-1",
operation_id="reindexDocuments",
bucket="import",
billing_cost=20,
rate_limit_cost=20,
)
]
assert billing.usages == rate_limit.usages
def test_direct_operation_admission_commits_at_successful_capability_issuance() -> None:
rate_limit = RecordingRateLimit()
charge = RecordingCharge()
broker = RecordingBroker()
admission = KnowledgeFSOperationAdmissionService(
rate_limit=rate_limit,
billing=RecordingBilling(charge),
)
metrics = MagicMock()
service = KnowledgeFSDirectOperationAdmissionService(admission=admission, broker=broker, metrics=metrics)
issued = service.issue_interactive(
tenant_id="tenant-1",
account_id="account-1",
control_space_id="control-1",
operation_id="createQuery",
)
assert issued.token == "interactive-capability"
assert charge.commits == 1
assert charge.refunds == 0
assert broker.interactive_calls == [
{
"tenant_id": "tenant-1",
"account_id": "account-1",
"control_space_id": "control-1",
"operation_id": "createQuery",
"resource_id": None,
"trace_id": None,
}
]
assert rate_limit.usages == [KnowledgeFSOperationUsage("tenant-1", "createQuery", "direct", 20, 20)]
assert [call.args[0] for call in metrics.record_operation_admission.call_args_list] == [
("createQuery", "direct", "reserve", "success"),
("createQuery", "direct", "commit", "success"),
]
assert "tenant-1" not in str(metrics.record_operation_admission.call_args_list)
def test_direct_operation_admission_refunds_when_service_capability_issuance_fails() -> None:
charge = RecordingCharge()
broker = RecordingBroker(error=PermissionError("credential lost access"))
admission = KnowledgeFSOperationAdmissionService(
rate_limit=RecordingRateLimit(),
billing=RecordingBilling(charge),
)
metrics = MagicMock()
service = KnowledgeFSDirectOperationAdmissionService(admission=admission, broker=broker, metrics=metrics)
profile = SimpleNamespace(tenant_id="tenant-1")
with pytest.raises(PermissionError, match="lost access"):
service.issue_service(profile=profile, operation_id="createQuery")
assert charge.commits == 0
assert charge.refunds == 1
assert broker.service_calls == [
{
"profile": profile,
"operation_id": "createQuery",
"resource_id": None,
"trace_id": None,
}
]
assert [call.args[0] for call in metrics.record_operation_admission.call_args_list] == [
("createQuery", "direct", "reserve", "success"),
("createQuery", "direct", "refund", "success"),
]
def test_direct_operation_admission_records_reserve_failure_without_tenant_labels() -> None:
admission = MagicMock()
admission.reserve.side_effect = KnowledgeFSOperationRateLimitExceededError("limited")
metrics = MagicMock()
service = KnowledgeFSDirectOperationAdmissionService(
admission=admission,
broker=RecordingBroker(),
metrics=metrics,
)
with pytest.raises(KnowledgeFSOperationRateLimitExceededError):
service.issue_interactive(
tenant_id="tenant-secret",
account_id="account-1",
control_space_id="control-1",
operation_id="createQuery",
)
assert metrics.record_operation_admission.call_args.args[0] == (
"createQuery",
"direct",
"reserve",
"failure",
)
assert "tenant-secret" not in str(metrics.record_operation_admission.call_args_list)
class FakeRedis:
def __init__(self, result: int) -> None:
self.result = result
self.calls: list[tuple[object, ...]] = []
def eval(self, *args: object) -> int:
self.calls.append(args)
return self.result
class RecordingRateLimitAudit:
def __init__(self) -> None:
self.usages: list[tuple[KnowledgeFSOperationUsage, str]] = []
def record_rejection(self, usage: KnowledgeFSOperationUsage, *, subscription_plan: str) -> None:
self.usages.append((usage, subscription_plan))
def test_weighted_rate_limit_uses_bucket_and_cost_atomically() -> None:
redis = FakeRedis(1)
audit = RecordingRateLimitAudit()
port = DifyKnowledgeFSWeightedRateLimitPort(
redis=redis,
audit=audit,
rate_limit_lookup=lambda _tenant_id: SimpleNamespace(enabled=True, limit=30, subscription_plan="sandbox"),
clock_ms=lambda: 1_000_000,
member_id=lambda: "request-1",
)
usage = KnowledgeFSOperationUsage("tenant-1", "createResearchTask", "query", 25, 25)
port.admit(usage)
assert len(redis.calls) == 1
_, key_count, key, now, window_start, limit, cost, member = redis.calls[0]
assert key_count == 1
assert key == "knowledge_fs:rate_limit:tenant-1:query"
assert (now, window_start, limit, cost, member) == (1_000_000, 940_000, 30, 25, "request-1")
assert audit.usages == []
def test_weighted_rate_limit_rejects_and_records_operation_without_external_io() -> None:
redis = FakeRedis(0)
audit = RecordingRateLimitAudit()
port = DifyKnowledgeFSWeightedRateLimitPort(
redis=redis,
audit=audit,
rate_limit_lookup=lambda _tenant_id: SimpleNamespace(enabled=True, limit=10, subscription_plan="sandbox"),
clock_ms=lambda: 1_000_000,
member_id=lambda: "request-1",
)
usage = KnowledgeFSOperationUsage("tenant-1", "reindexDocuments", "import", 20, 20)
with pytest.raises(KnowledgeFSOperationRateLimitExceededError):
port.admit(usage)
assert audit.usages == [(usage, "sandbox")]
class FakeBillingGateway:
def __init__(self, reservation_id: str | None = "reservation-1", *, fail: bool = False) -> None:
self.reservation_id = reservation_id
self.fail = fail
self.calls: list[tuple[str, dict[str, object]]] = []
def quota_reserve(self, **kwargs: object) -> dict[str, str]:
self.calls.append(("reserve", kwargs))
if self.fail:
raise RuntimeError("billing unavailable")
return {"reservation_id": self.reservation_id} if self.reservation_id else {}
def quota_commit(self, **kwargs: object) -> dict[str, str]:
self.calls.append(("commit", kwargs))
return {"result": "success"}
def quota_release(self, **kwargs: object) -> dict[str, str]:
self.calls.append(("release", kwargs))
return {"result": "success"}
def test_billing_reservation_commits_or_releases_exact_operation_cost() -> None:
gateway = FakeBillingGateway()
port = DifyKnowledgeFSBillingPort(
gateway=gateway,
billing_enabled=lambda: True,
request_id=lambda: "request-1",
)
usage = KnowledgeFSOperationUsage("tenant-1", "createResearchTask", "query", 25, 25)
committed = port.reserve(usage)
committed.commit()
committed.commit()
released = port.reserve(usage)
released.refund()
released.refund()
assert [kind for kind, _ in gateway.calls] == ["reserve", "commit", "reserve", "release"]
assert gateway.calls[0][1] == {
"amount": 25,
"bucket": "query",
"feature_key": "knowledge_fs_operations",
"meta": {"operation_id": "createResearchTask", "source": "knowledge_fs"},
"request_id": "request-1",
"tenant_id": "tenant-1",
}
assert gateway.calls[1][1]["actual_amount"] == 25
assert gateway.calls[1][1]["reservation_id"] == "reservation-1"
assert gateway.calls[3][1]["reservation_id"] == "reservation-1"
def test_billing_explicit_exhaustion_fails_closed_but_transport_failure_fails_open() -> None:
usage = KnowledgeFSOperationUsage("tenant-1", "getSettings", "read", 1, 1)
exhausted = DifyKnowledgeFSBillingPort(
gateway=FakeBillingGateway(reservation_id=None),
billing_enabled=lambda: True,
)
unavailable = DifyKnowledgeFSBillingPort(
gateway=FakeBillingGateway(fail=True),
billing_enabled=lambda: True,
)
with pytest.raises(KnowledgeFSOperationQuotaExceededError):
exhausted.reserve(usage)
charge = unavailable.reserve(usage)
charge.commit()
charge.refund()

View File

@ -122,18 +122,12 @@ def test_ready_product_operations_exactly_match_capability_method_path_and_actio
assert product_operation.action == capability_operation.action
assert product_operation.resource_resolver == capability_operation.resource_type
assert product_operation.permission == product_operation.rbac_permission
assert product_operation.billing_cost > 0
assert product_operation.rate_limit_cost > 0
assert product_operation.rate_limit_bucket in {"direct", "import", "job", "query", "read", "write"}
assert product_operation.max_request_bytes >= 0
assert product_operation.max_response_bytes >= 0
if product_operation.transport == "json":
assert product_operation.stream_kind == "json"
assert product_operation.max_response_bytes > 0
assert KNOWLEDGE_FS_PRODUCT_OPERATIONS["reindexDocuments"].rate_limit_bucket == "import"
assert KNOWLEDGE_FS_PRODUCT_OPERATIONS["createResearchTask"].rate_limit_bucket == "query"
assert KNOWLEDGE_FS_PRODUCT_OPERATIONS["getCompilationJob"].rate_limit_bucket == "job"
assert KNOWLEDGE_FS_PRODUCT_OPERATIONS["createDocument"].max_request_bytes == 15 * 1024 * 1024
assert KNOWLEDGE_FS_PRODUCT_OPERATIONS["importSourceWorkflow"].max_request_bytes == 4 * 1024 * 1024
assert KNOWLEDGE_FS_PRODUCT_OPERATIONS["uploadSmallFile"].max_request_bytes == 8 * 1024 * 1024

View File

@ -22,9 +22,7 @@ def test_runtime_wires_one_shared_authorization_and_remote_graph(monkeypatch: py
monkeypatch.setattr(runtime.dify_config, "KNOWLEDGE_FS_PRODUCT_MAX_RESPONSE_BYTES", 4096)
factory_names = (
"DifyKnowledgeFSBillingPort",
"DifyKnowledgeFSProductRBACPort",
"DifyKnowledgeFSWeightedRateLimitPort",
"HTTPKnowledgeFSProductRemoteClient",
"KnowledgeFSAppAdmissionService",
"KnowledgeFSAppBindingManagementService",
@ -35,14 +33,11 @@ def test_runtime_wires_one_shared_authorization_and_remote_graph(monkeypatch: py
"KnowledgeFSControlSpaceCommandService",
"KnowledgeFSCredentialService",
"KnowledgeFSDataFacade",
"KnowledgeFSDirectOperationAdmissionService",
"KnowledgeFSOperationAdmissionService",
"KnowledgeFSProductApplicationService",
"KnowledgeFSProductService",
"KnowledgeFSRevocationCommandProducer",
"KnowledgeFSWorkspaceCutoverService",
"KnowledgeFSWorkspaceGreenfieldInitializer",
"LoggingKnowledgeFSRateLimitAudit",
"SQLKnowledgeFSAppCatalog",
"SQLKnowledgeFSWorkspaceMemberPort",
"SQLKnowledgeFSWorkspaceRuntimeGate",
@ -67,9 +62,7 @@ def test_runtime_wires_one_shared_authorization_and_remote_graph(monkeypatch: py
assert result.broker is factories["KnowledgeFSCapabilityBroker"].return_value
assert result.control_plane is factories["KnowledgeFSControlPlaneService"].return_value
assert result.credentials is factories["KnowledgeFSCredentialService"].return_value
assert result.direct_operation_admission is factories["KnowledgeFSDirectOperationAdmissionService"].return_value
assert result.facade is factories["KnowledgeFSDataFacade"].return_value
assert result.operation_admission is factories["KnowledgeFSOperationAdmissionService"].return_value
factories["HTTPKnowledgeFSProductRemoteClient"].assert_called_once_with(
base_url="https://knowledge-fs.test",