Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
108 changes: 19 additions & 89 deletions backend/ext_components/aidp/apps/aidp_mgmt_app.py
Original file line number Diff line number Diff line change
Expand Up @@ -18,25 +18,27 @@
import time
from http import HTTPStatus
from typing import Annotated, List, Optional
from uuid import UUID

from fastapi import APIRouter, File, HTTPException, Path, Query, Request, UploadFile
from fastapi.responses import JSONResponse, StreamingResponse
from nexent.core.concurrency import run_blocking
from fastapi import APIRouter, File, Path, Query, Request, UploadFile
from fastapi.responses import JSONResponse
from pydantic import BaseModel, Field
from sqlalchemy.exc import IntegrityError
from starlette.background import BackgroundTask
from nexent.core.concurrency import run_blocking

from consts.const import AIDP_API_KEY, AIDP_SERVER_URL
from consts.error_code import ErrorCode
from consts.exceptions import AppException, UnauthorizedError
from database.user_tenant_db import get_user_role_by_tenant
from ext_components.aidp.consts.aidp_exceptions import (
AidpGroupValidationError,
AidpKbConflictError,
AidpKbNotFoundError,
AidpKbPermissionDeniedError,
AidpKbSyncError,
AidpGroupValidationError,
)
from ext_components.aidp.database import aidp_permission_db
from ext_components.aidp.services import aidp_permission_service as perms
from ext_components.aidp.services.aidp_kb_update_service import save_kb_settings
from ext_components.aidp.services.aidp_access_service import (
get_cached_aidp_channels,
get_cached_aidp_doc_count,
Expand All @@ -46,13 +48,6 @@
invalidate_aidp_kb_detail_cache,
resolve_current_aidp_access,
)
from ext_components.aidp.services.aidp_kb_update_service import save_kb_settings
from ext_components.aidp.services.aidp_permission_service import (
EDIT,
PRIVATE,
READ_ONLY,
_validate_group_ids_strict, # noqa: F401 - retained as a module-level compatibility symbol
)
from ext_components.aidp.services.aidp_service import (
_timestamp_to_iso,
count_aidp_docs_impl,
Expand All @@ -63,15 +58,18 @@
list_aidp_doc_history_impl,
list_aidp_docs_impl,
list_aidp_models_impl,
remove_aidp_docs_impl,
stream_aidp_doc_impl,
select_aidp_channel,
update_aidp_kb_impl,
upload_aidp_docs_impl,
)
from ext_components.aidp.services.aidp_permission_service import (
EDIT,
PRIVATE,
READ_ONLY,
_validate_group_ids_strict,
)
from utils import auth_utils as auth_utils_module


aidp_mgmt_router = APIRouter(prefix="/aidp-mgmt")
logger = logging.getLogger("aidp_mgmt_app")

Expand Down Expand Up @@ -201,17 +199,6 @@
)


class RemoveAidpDocumentsRequest(BaseModel):

file_uuids: List[UUID] = Field(..., min_length=1, description="AIDP file UUIDs")


class DownloadAidpDocumentRequest(BaseModel):
"""AIDP file selected for download."""

file_uuid: UUID = Field(..., description="AIDP file UUID")


# ---------------------------------------------------------------------------
# Auth helpers
# ---------------------------------------------------------------------------
Expand Down Expand Up @@ -268,6 +255,11 @@
)


# HTTPException is imported lazily to keep FastAPI's exception handler in
# control of the response body.
from fastapi import HTTPException # noqa: E402 (placed here to avoid editing mid-file)


def _credentials() -> tuple[str, str]:
return AIDP_SERVER_URL, AIDP_API_KEY

Expand Down Expand Up @@ -367,13 +359,13 @@
)
return
_HISTORY_FALLBACK_REPORTED.add(kds_id)
logger.warning(
"AIDP all-status file history unavailable for KB %s (%s); the document list falls "
"back to ingested files only, so files under processing stay invisible and the "
"status column stays empty",
kds_id,
reason,
)

Check warning on line 368 in backend/ext_components/aidp/apps/aidp_mgmt_app.py

View check run for this annotation

SonarQubeCloud / SonarCloud Code Analysis

Log Injection via unsanitized user input in aidp_mgmt_app._log_history_fallback()

See more on https://sonarcloud.io/project/issues?id=ModelEngine-Group_nexent&issues=AaDCHPy1e8_wdapXFZZb&open=AaDCHPy1e8_wdapXFZZb&pullRequest=3983


def _resolve_doc_history_channel(
Expand Down Expand Up @@ -549,9 +541,6 @@
server_url, api_key = _credentials()
normalized_keyword = (keyword or "").strip() or None
started_at = time.perf_counter()
# ``run_blocking`` (not a bare to_thread) keeps AIDP catalog reads on the
# managed control-io lane, and the keyword is still forwarded so the remote
# catalog is narrowed server-side before the permission intersection.
rows = await run_blocking(
"aidp-accessible-rows",
_current_accessible_rows,
Expand Down Expand Up @@ -947,16 +936,16 @@
history_result = await _load_doc_history(server_url, api_key, kds_id)
if history_result is not None:
result = _paginate_history_documents(history_result, page, page_size)
logger.info(
"AIDP document list timing: total_ms=%.1f kb_id=%s page=%d page_size=%d "
"page_count=%d total_count=%d total_reliable=True source=history",
(time.perf_counter() - started_at) * 1000,
kds_id,
page,
page_size,
len(result["value"]),
result["total_count"],
)

Check warning on line 948 in backend/ext_components/aidp/apps/aidp_mgmt_app.py

View check run for this annotation

SonarQubeCloud / SonarCloud Code Analysis

Log Injection via unsanitized user input in aidp_mgmt_app.list_documents()

See more on https://sonarcloud.io/project/issues?id=ModelEngine-Group_nexent&issues=AaDCHPy1e8_wdapXFZZc&open=AaDCHPy1e8_wdapXFZZc&pullRequest=3983
return JSONResponse(status_code=HTTPStatus.OK, content=result)

# Fallback: the completed-files listing (historical behaviour), used when
Expand Down Expand Up @@ -1029,65 +1018,6 @@
return JSONResponse(status_code=HTTPStatus.OK, content=result)


@aidp_mgmt_router.post("/knowledge-bases/{kds_id}/documents/remove")
async def remove_documents(
request: Request,
kds_id: Annotated[str, Path(description="Knowledge base ID")],
body: RemoveAidpDocumentsRequest,
) -> JSONResponse:
"""Remove AIDP documents."""
user_id, tenant_id = await _auth(request)
perms.require_permission(kds_id, user_id, tenant_id, required="EDIT")

server_url, api_key = _credentials()
result = await run_blocking(
"aidp-remove-documents",
remove_aidp_docs_impl,
server_url,
api_key,
kds_id,
[str(file_uuid) for file_uuid in body.file_uuids],
lane="control-io",
owner="config",
)

success_list = result["success_list"]
if success_list:
invalidate_aidp_kb_detail_cache(server_url, api_key, kds_id)
invalidate_aidp_doc_count_cache(server_url, api_key, kds_id)
return JSONResponse(status_code=HTTPStatus.OK, content=result)


@aidp_mgmt_router.post("/knowledge-bases/{kds_id}/documents/download")
async def download_document(
request: Request,
kds_id: Annotated[str, Path(description="Knowledge base ID")],
body: DownloadAidpDocumentRequest,
) -> StreamingResponse:
"""Proxy an AIDP document as a binary attachment."""
user_id, tenant_id = await _auth(request)
perms.require_permission(kds_id, user_id, tenant_id, required="READ")

server_url, api_key = _credentials()
aidp_response = await stream_aidp_doc_impl(
server_url,
api_key,
kds_id,
str(body.file_uuid),
)
response_headers = {
"Content-Disposition": aidp_response.headers["Content-Disposition"],
"X-File-Size": aidp_response.headers["X-File-Size"],
}

return StreamingResponse(
aidp_response.aiter_bytes(),
media_type=aidp_response.headers["Content-Type"],
headers=response_headers,
background=BackgroundTask(aidp_response.aclose),
)


@aidp_mgmt_router.patch("/aidp-permissions/{kds_id}")
async def set_permission(
request: Request,
Expand Down
140 changes: 2 additions & 138 deletions backend/ext_components/aidp/services/aidp_service.py
Original file line number Diff line number Diff line change
Expand Up @@ -5,16 +5,15 @@
import logging
import time
from datetime import datetime, timezone
from typing import Any, Callable, Dict, List, NoReturn
from typing import Any, Callable, Dict, List
from urllib.parse import quote, urljoin

import httpx
from nexent.utils.http_client_manager import http_client_manager

from consts.const import AIDP_TENANT_ID
from consts.error_code import ErrorCode
from consts.exceptions import AppException

from nexent.utils.http_client_manager import http_client_manager

logger = logging.getLogger("aidp_service")

Expand Down Expand Up @@ -67,39 +66,6 @@
return reason[:_MAX_UPSTREAM_ERROR_REASON_LENGTH]


def _raise_aidp_http_error(error: httpx.HTTPStatusError, operation: str) -> NoReturn:
"""Map an AIDP HTTP error to the common application exception format."""
response = error.response
upstream_reason = _extract_upstream_error(response)
logger.exception(
"AIDP %s HTTP error: status_code=%s upstream_reason=%s",
operation,
response.status_code,
upstream_reason or "unavailable",
)
details = {
"upstream_status": response.status_code,
"upstream_reason": upstream_reason,
}
error_code = {
401: ErrorCode.AIDP_AUTH_ERROR,
403: ErrorCode.AIDP_AUTH_ERROR,
429: ErrorCode.AIDP_RATE_LIMIT,
}.get(response.status_code, ErrorCode.AIDP_SERVICE_ERROR)
fallback_message = {
ErrorCode.AIDP_AUTH_ERROR: f"AIDP authentication failed: {str(error)}",
ErrorCode.AIDP_RATE_LIMIT: f"AIDP rate limit exceeded: {str(error)}",
}.get(
error_code,
f"AIDP API HTTP error {response.status_code}: {str(error)}",
)
raise AppException(
error_code,
upstream_reason or fallback_message,
details=details,
)


def _extract_upload_failures(response: httpx.Response) -> List[Dict[str, str]]:
"""Extract per-file upload failures from AIDP's structured error body."""
try:
Expand Down Expand Up @@ -364,7 +330,6 @@
_AIDP_RETRY_BACKOFF_FACTOR = 0.5
_AIDP_RETRYABLE_STATUS_CODES = {408, 429, 500, 502, 503, 504}
_AIDP_READ_TIMEOUT_SECONDS = 30.0
_AIDP_DOWNLOAD_TIMEOUT_SECONDS = 120.0


def _request_with_retry(
Expand Down Expand Up @@ -1249,107 +1214,6 @@
)


def remove_aidp_docs_impl(
server_url: str,
api_key: str,
kds_id: str,
file_uuids: List[str],
) -> Dict[str, Any]:
"""Remove one or more documents from an AIDP knowledge base."""
normalized_url = _validate_params(server_url, api_key)

headers = {
"Authorization": f"Bearer {api_key}",
"Content-Type": "application/json",
}
remove_path = f"{_get_list_path()}/{kds_id}/KnowledgeFiles/Remove"
remove_url = urljoin(f"{normalized_url}/", remove_path)
logger.info("Removing %d AIDP documents from %s", len(file_uuids), remove_url)

try:
client = http_client_manager.get_sync_client(
base_url=normalized_url,
timeout=_AIDP_READ_TIMEOUT_SECONDS,
verify_ssl=False,
)
response = _request_with_retry(
lambda: client.post(
remove_url,
headers=headers,
json={"file_uuids": file_uuids},
),
context=f"remove-docs:{kds_id}",
)
response.raise_for_status()
result = response.json()
return result
except httpx.RequestError as e:
logger.exception("AIDP document removal request failed: %s", e)
raise AppException(
ErrorCode.AIDP_CONNECTION_ERROR,
f"AIDP API request failed: {str(e)}",
)
except httpx.HTTPStatusError as e:
_raise_aidp_http_error(e, "document removal")
except ValueError as e:
logger.exception("Failed to parse AIDP document removal response: %s", e)
raise AppException(
ErrorCode.AIDP_RESPONSE_ERROR,
f"Failed to parse AIDP API response: {str(e)}",
)


async def stream_aidp_doc_impl(
server_url: str,
api_key: str,
kds_id: str,
file_uuid: str,
) -> httpx.Response:
"""Open a streaming response for one AIDP document."""
normalized_url = _validate_params(server_url, api_key)

headers = {
"Authorization": f"Bearer {api_key}",
"Content-Type": "application/json",
}
download_path = f"{_get_list_path()}/{kds_id}/KnowledgeFiles/Download"
download_url = urljoin(f"{normalized_url}/", download_path)
logger.info("Downloading AIDP document %s from %s", file_uuid, download_url)

response: httpx.Response | None = None
try:
client = http_client_manager.get_async_client(
base_url=normalized_url,
timeout=_AIDP_DOWNLOAD_TIMEOUT_SECONDS,
verify_ssl=False,
)
response = await client.send(
client.build_request(
"POST",
download_url,
headers=headers,
json={"file_uuid": file_uuid},
),
stream=True,
)
if response.status_code >= 400:
await response.aread()
response.raise_for_status()
return response
except httpx.RequestError as e:
if response is not None:
await response.aclose()
logger.exception("AIDP document download request failed: %s", e)
raise AppException(
ErrorCode.AIDP_CONNECTION_ERROR,
f"AIDP API request failed: {str(e)}",
)
except httpx.HTTPStatusError as e:
if response is not None:
await response.aclose()
_raise_aidp_http_error(e, "document download")


def count_aidp_docs_impl(server_url: str, api_key: str, kds_id: str) -> int:
"""Get total document count in a KB via AIDP POST .../Count endpoint.

Expand Down Expand Up @@ -1725,7 +1589,7 @@
)


def list_aidp_doc_history_impl(

Check failure on line 1592 in backend/ext_components/aidp/services/aidp_service.py

View check run for this annotation

SonarQubeCloud / SonarCloud Code Analysis

Refactor this function to reduce its Cognitive Complexity from 16 to the 15 allowed.

See more on https://sonarcloud.io/project/issues?id=ModelEngine-Group_nexent&issues=AaDCHP0Ge8_wdapXFZZd&open=AaDCHP0Ge8_wdapXFZZd&pullRequest=3983
server_url: str,
api_key: str,
fs_id: str,
Expand Down
Loading
Loading