diff --git a/backend/ext_components/aidp/apps/aidp_mgmt_app.py b/backend/ext_components/aidp/apps/aidp_mgmt_app.py index 5fa9fe224..23e768419 100644 --- a/backend/ext_components/aidp/apps/aidp_mgmt_app.py +++ b/backend/ext_components/aidp/apps/aidp_mgmt_app.py @@ -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, @@ -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, @@ -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") @@ -201,17 +199,6 @@ class SetPermissionRequest(BaseModel): ) -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 # --------------------------------------------------------------------------- @@ -268,6 +255,11 @@ def _raise_aidp_conflict(exc: IntegrityError) -> None: ) +# 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 @@ -549,9 +541,6 @@ async def list_knowledge_bases( 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, @@ -1029,65 +1018,6 @@ async def list_documents( 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, diff --git a/backend/ext_components/aidp/services/aidp_service.py b/backend/ext_components/aidp/services/aidp_service.py index 50400aee4..4116e5fee 100644 --- a/backend/ext_components/aidp/services/aidp_service.py +++ b/backend/ext_components/aidp/services/aidp_service.py @@ -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") @@ -67,39 +66,6 @@ def _extract_upstream_error(response: httpx.Response) -> str | None: 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: @@ -364,7 +330,6 @@ def _validate_params(server_url: str, api_key: str) -> str: _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( @@ -1249,107 +1214,6 @@ def upload_aidp_docs_impl( ) -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. diff --git a/backend/services/human_interaction/application.py b/backend/services/human_interaction/application.py index d2251fa32..bf84b5e71 100644 --- a/backend/services/human_interaction/application.py +++ b/backend/services/human_interaction/application.py @@ -112,7 +112,24 @@ def authorize(): except Exception as exc: raise RunTerminated("Run authorization could not be revalidated") from exc - port = RuntimeInteractionPort(service, identity, lease.owner_id, authorize, live_resume=True) + def report_waiting(waiting: bool) -> None: + # Worker threads park inside _wait_until_ready; relay the state to + # the scheduler loop so human waits do not consume execution slots. + # The reporter is a pure concurrency-budget optimization: a dead + # loop (service shutdown) or a stale SDK copy without mark_waiting + # must degrade to slot-consuming waits, never break execution. + mark_waiting = getattr(human_run_scheduler, "mark_waiting", None) + if mark_waiting is None: + return + # call_soon_threadsafe takes no kwargs: wrap the keyword-only flag. + try: + loop.call_soon_threadsafe(lambda: mark_waiting(lease.job_id, waiting=waiting)) + except RuntimeError: + pass + + port = RuntimeInteractionPort( + service, identity, lease.owner_id, authorize, live_resume=True, wait_reporter=report_waiting, + ) saved = port.request_payload if saved.get("runtime_mode") == "native-live-v1" and port.checkpoint: raise RecoveryRequired("The original native execution is no longer available") @@ -173,24 +190,21 @@ async def _flush_if_due() -> None: """Flush buffered chunks on the interval/batch trigger to keep DB order. The async loop flushes on its own; the worker flushes again before - each HITL transaction. Peek first and never take-put-back — that - would open a transiently empty window the worker's idle poll can - mistake for "async is done". + each HITL transaction. """ nonlocal last_flush - if (time.monotonic() - last_flush >= _FLUSH_INTERVAL - or port.peek_chunks() >= _FLUSH_BATCH): - buffered = port.take_chunks() - if buffered: - try: - port.begin_emit() - await run_blocking( - "hitl-port-emit_chunks", port.emit_chunks, buffered, - lane="control-io", owner=__name__, - ) - last_flush = time.monotonic() - finally: - port.end_emit() + if ((time.monotonic() - last_flush >= _FLUSH_INTERVAL + or port.peek_chunks() >= _FLUSH_BATCH) + and port.peek_chunks()): + # drain_and_emit runs take_chunks + emit_chunks in one + # emit-lock critical section on the lane thread, so it cannot + # interleave with the worker's flush — seq order stays + # production order. + await run_blocking( + "hitl-port-drain_and_emit", port.drain_and_emit, + lane="control-io", owner=__name__, + ) + last_flush = time.monotonic() chunk_iter = _stream_agent_chunks( agent_request=request, user_id=identity["user_id"], tenant_id=identity["tenant_id"], @@ -216,6 +230,10 @@ async def _flush_if_due() -> None: try: chunk = task.result() except StopAsyncIteration: + # Stream exhausted: drop the consumed task so the finally + # block does not re-raise its StopAsyncIteration and skip + # the leftover flush plus the terminal finish() write. + anext_task = None break port.add_chunk(chunk) await _flush_if_due() @@ -230,16 +248,10 @@ async def _flush_if_due() -> None: except Exception: pass # Final flush so leftover chunks precede finish() in DB order. - leftover = port.take_chunks() - if leftover: - try: - port.begin_emit() - await run_blocking( - "hitl-port-emit_chunks", port.emit_chunks, leftover, - lane="control-io", owner=__name__, - ) - finally: - port.end_emit() + await run_blocking( + "hitl-port-drain_and_emit", port.drain_and_emit, + lane="control-io", owner=__name__, + ) await run_blocking( "hitl-port-finish", port.finish, run_info.attempt_outcome or "failed", lane="control-io", owner=__name__, diff --git a/backend/services/human_interaction/runtime_port.py b/backend/services/human_interaction/runtime_port.py index 62b01b417..d0345fb54 100644 --- a/backend/services/human_interaction/runtime_port.py +++ b/backend/services/human_interaction/runtime_port.py @@ -17,14 +17,18 @@ class RuntimeInteractionPort: - def __init__(self, service, identity, owner_id, authorize, allowed_tools=(), *, live_resume=False, stop_event=None): + def __init__(self, service, identity, owner_id, authorize, allowed_tools=(), *, live_resume=False, stop_event=None, + wait_reporter=None): """Bind the port and set up async-worker chunk coordination. ``_chunk_buffer`` stages processed observer chunks: the async consumer appends via add_chunk and the worker flushes before each HITL event so chunk rows precede human_interaction rows in DB event order. - ``_emit_in_flight`` marks an in-flight hand-off to the DB thread so an - empty buffer is not mistaken for "async is idle". + ``_emit_lock`` serializes every take_chunks + emit_chunks critical + section (async drain_and_emit and the worker's own flush), so a held + lock also means "a drain is in flight and its chunks are uncommitted". + ``wait_reporter`` relays enter/exit of human-input waits to the + scheduler so parked runs stop consuming execution concurrency slots. """ self.service = service self.repository = service.repository @@ -38,9 +42,10 @@ def __init__(self, service, identity, owner_id, authorize, allowed_tools=(), *, self.allowed_tools = frozenset(allowed_tools) self.live_resume = live_resume self.stop_event = stop_event + self.wait_reporter = wait_reporter self._chunk_buffer: list[str] = [] self._chunk_buffer_lock = threading.Lock() - self._emit_in_flight = threading.Event() + self._emit_lock = threading.Lock() with self.transaction() as tx: self.checkpoint = self.cipher.open(tx.run.checkpoint) self.request_payload = self.cipher.open(tx.run.request_payload) @@ -62,48 +67,59 @@ def peek_chunks(self) -> int: with self._chunk_buffer_lock: return len(self._chunk_buffer) - def begin_emit(self) -> None: - """Signal that the async consumer is handing chunks to run_blocking(emit_chunks).""" - self._emit_in_flight.set() + def drain_and_emit(self) -> None: + """Atomically drain and persist buffered chunks. Lane-thread only. - def end_emit(self) -> None: - """Signal that the async consumer's run_blocking(emit_chunks) has returned.""" - self._emit_in_flight.clear() + take_chunks and emit_chunks share one critical section under the emit + lock so a concurrent flush can never interleave DB seq assignment — + the seq order then always matches the chunk production order. + """ + with self._emit_lock: + chunks = self.take_chunks() + if not chunks: + return + try: + self.emit_chunks(chunks) + except Exception: + # Never lose drained chunks; put them back for retry. + for chunk in chunks: + self.add_chunk(chunk) + raise def flush_chunks_until_idle(self, *, max_wait_ms: int = 500, settle_ms: int = 20) -> None: """Wait for the async consumer to drain the observer queue, then persist. - Worker-thread only. Returns once the buffer has stayed empty for - ``settle_ms`` with no in-flight emit; ``max_wait_ms`` is a hard - upper bound measured from entry and never resets. + Worker-thread only. take_chunks and emit_chunks run inside the same + emit-lock critical section as the async drain_and_emit, so concurrent + flushes cannot reorder seq assignment. ``max_wait_ms`` bounds how long + we wait for NEW chunks — it must NOT fire while another drain holds + the lock, because its chunks are produced before this HITL event but + still uncommitted; returning early would order HITL rows ahead of + them. The deadline may therefore only fire once the lock is free. """ hard_deadline = time.monotonic() + max_wait_ms / 1000.0 idle_since: float | None = None while True: now = time.monotonic() - if now >= hard_deadline: + if now >= hard_deadline and not self._emit_lock.locked(): break - - chunks = self.take_chunks() - if chunks: - try: - self.emit_chunks(chunks) - except Exception: - # Never lose drained chunks; put them back for retry. - for chunk in chunks: - self.add_chunk(chunk) - raise - idle_since = None - elif self._emit_in_flight.is_set(): - # Empty buffer is not idle while an emit is in flight. - idle_since = None - elif idle_since is None: - idle_since = now - elif now - idle_since >= settle_ms / 1000.0: - break - - time.sleep(min(settle_ms / 1000.0, hard_deadline - now)) + with self._emit_lock: + chunks = self.take_chunks() + if chunks: + try: + self.emit_chunks(chunks) + except Exception: + # Never lose drained chunks; put them back for retry. + for chunk in chunks: + self.add_chunk(chunk) + raise + idle_since = None + elif idle_since is None: + idle_since = now + elif now - idle_since >= settle_ms / 1000.0: + break + time.sleep(min(settle_ms / 1000.0, max(hard_deadline - now, 0.0))) @contextmanager def transaction(self, *, receipt=False): @@ -199,33 +215,39 @@ def close_steering(self, checkpoint): def _wait_until_ready(self): """Park this worker while retaining the current Python continuation.""" - next_authorization_check = 0.0 - while True: - if self.stop_event is not None and self.stop_event.is_set(): - raise RunTerminated("The managed execution was cancelled") - now = time.monotonic() - if now >= next_authorization_check: - self.authorize() - next_authorization_check = now + 5.0 - with self.repository.transaction(self.run_id, self.tenant_id, self.user_id) as tx: - if (tx is None or tx.run.fence != self.fence or tx.run.lock_owner != self.owner_id - or tx.run.lock_until is None or tx.run.lock_until <= utcnow()): - raise RunTerminated("Execution lease is no longer valid") - self.service._expire(tx) - if tx.run.status == "READY": - # Flush so lingering chunks precede the human_run row. - self.flush_chunks_until_idle() - tx.run.status = "RUNNING" - tx.emit({"type": "human_run", "content": { - "run_id": self.run_id, "status": "RUNNING", - }}) - return - if tx.run.status != "WAITING_HUMAN": - raise RunTerminated("Run no longer permits live continuation") - if self.stop_event is not None: - self.stop_event.wait(0.2) - else: - time.sleep(0.2) + if self.wait_reporter is not None: + self.wait_reporter(True) + try: + next_authorization_check = 0.0 + while True: + if self.stop_event is not None and self.stop_event.is_set(): + raise RunTerminated("The managed execution was cancelled") + now = time.monotonic() + if now >= next_authorization_check: + self.authorize() + next_authorization_check = now + 5.0 + with self.repository.transaction(self.run_id, self.tenant_id, self.user_id) as tx: + if (tx is None or tx.run.fence != self.fence or tx.run.lock_owner != self.owner_id + or tx.run.lock_until is None or tx.run.lock_until <= utcnow()): + raise RunTerminated("Execution lease is no longer valid") + self.service._expire(tx) + if tx.run.status == "READY": + # Flush so lingering chunks precede the human_run row. + self.flush_chunks_until_idle() + tx.run.status = "RUNNING" + tx.emit({"type": "human_run", "content": { + "run_id": self.run_id, "status": "RUNNING", + }}) + return + if tx.run.status != "WAITING_HUMAN": + raise RunTerminated("Run no longer permits live continuation") + if self.stop_event is not None: + self.stop_event.wait(0.2) + else: + time.sleep(0.2) + finally: + if self.wait_reporter is not None: + self.wait_reporter(False) def dispatch(self, slot, tool, arguments, *, interaction=None): self.authorize() diff --git a/deploy/env/.env.example b/deploy/env/.env.example index 7c3b478fe..58f7d817e 100644 --- a/deploy/env/.env.example +++ b/deploy/env/.env.example @@ -360,5 +360,5 @@ HITL_ACCEPT_NEW_RUNS=true # non-built-in Agent tool requires explicit review before execution. HITL_TOOL_APPROVAL_ENABLED=false HITL_ENCRYPTION_KEY= -HITL_WAIT_SECONDS=86400 -HITL_MAX_CONCURRENCY=2 +HITL_WAIT_SECONDS=3600 +HITL_MAX_CONCURRENCY=100 diff --git a/deploy/env/hitl.env.example b/deploy/env/hitl.env.example index 1f4d8b2d6..c72e58463 100644 --- a/deploy/env/hitl.env.example +++ b/deploy/env/hitl.env.example @@ -4,5 +4,5 @@ HITL_ENABLED=true HITL_ACCEPT_NEW_RUNS=true HITL_TOOL_APPROVAL_ENABLED=false HITL_ENCRYPTION_KEY= -HITL_WAIT_SECONDS=86400 -HITL_MAX_CONCURRENCY=2 +HITL_WAIT_SECONDS=3600 +HITL_MAX_CONCURRENCY=100 diff --git a/frontend/app/[locale]/newchat/adapter/remote-chat-model-adapter.ts b/frontend/app/[locale]/newchat/adapter/remote-chat-model-adapter.ts index f394ec5a0..ad973bad3 100644 --- a/frontend/app/[locale]/newchat/adapter/remote-chat-model-adapter.ts +++ b/frontend/app/[locale]/newchat/adapter/remote-chat-model-adapter.ts @@ -2191,10 +2191,21 @@ export const remoteChatModelAdapter: ChatModelAdapter = { } if (chunk.type === "human_run") { - const value = - typeof chunk.content === "string" + // The stream loop below has a finally but no catch: a malformed + // payload must not kill the whole chat stream. Skip the chunk and + // let the HITL controller snapshot/polling recover the state. + let value: Record; + try { + value = (typeof chunk.content === "string" ? JSON.parse(chunk.content) - : chunk.content; + : chunk.content) as Record; + } catch (error) { + log.warn( + "[ChatModelAdapter] Failed to parse human_run chunk:", + error + ); + continue; + } if (value && typeof value.run_id === "string") humanRunId = value.run_id; custom?.onHumanInteractionEvent?.(); diff --git a/frontend/app/[locale]/newchat/page.tsx b/frontend/app/[locale]/newchat/page.tsx index 31f40986e..0ac749d13 100644 --- a/frontend/app/[locale]/newchat/page.tsx +++ b/frontend/app/[locale]/newchat/page.tsx @@ -583,7 +583,9 @@ const HomeContent: FC<{ onRuntimeMetadataSent: handleRuntimeMetadataSent, onKnowledgeScopeResolved: handleKnowledgeScopeResolved, onGenerationStopped: handleGenerationStopped, - onHumanInteractionEvent: hitlController.refresh, + // HITL events (ask_user suspend, run transitions) must bypass the + // snapshot throttle: the event stream goes quiet afterwards. + onHumanInteractionEvent: () => hitlController.refresh(true), enablePlan: chatMode === "planning", enableHitl, ...(activeThreadId diff --git a/frontend/ext_components/aidp/components/AidpCreateKbModal.tsx b/frontend/ext_components/aidp/components/AidpCreateKbModal.tsx index 5e4eb93ed..53d94f68a 100644 --- a/frontend/ext_components/aidp/components/AidpCreateKbModal.tsx +++ b/frontend/ext_components/aidp/components/AidpCreateKbModal.tsx @@ -1,55 +1,70 @@ "use client"; -import React, { useEffect, useMemo, useRef, useState } from "react"; +import React, { useState, useMemo, useEffect, useRef } from "react"; import { useTranslation } from "react-i18next"; import { useQuery } from "@tanstack/react-query"; -import type { TFunction } from "i18next"; import { - Collapse, + Modal, Form, + Input, InputNumber, - Modal, - Select, + Steps, + Upload, + Button, + message, Space, + Divider, + Collapse, Switch, Tooltip, - Upload, - message, + Select, } from "antd"; -import { - InboxOutlined, - QuestionCircleOutlined, - SettingOutlined, -} from "@ant-design/icons"; +import { InboxOutlined, QuestionCircleOutlined } from "@ant-design/icons"; import type { AidpKnowledgeBaseItem } from "@/types/agentConfig"; -import type { - AidpModelItem, - AidpUploadResponse, -} from "@/ext_components/aidp/services/aidpKnowledgeService"; +import type { AidpModelItem } from "@/ext_components/aidp/services/aidpKnowledgeService"; import aidpKnowledgeService from "@/ext_components/aidp/services/aidpKnowledgeService"; -import { AIDP_ACCEPT_STRING } from "@/const/knowledgeBase"; +import { USER_ROLES } from "@/const/auth"; + +/** + * Antd's Upload component (Dragger) requires ``originFileObj`` to satisfy + * the ``RcFile`` shape (``File`` + ``uid`` + ``lastModifiedDate``). We + * store raw ``File`` objects in component state, so we cast at the + * render boundary. The structural cast is sufficient because antd does + * not read the extra fields — it only requires them to exist for type + * compatibility. + */ +type RcFileLike = File & { uid: string; lastModifiedDate: Date }; +import { + AIDP_ACCEPT_STRING, + AIDP_KNOWLEDGE_BASE_NAME_PATTERN, +} from "@/const/knowledgeBase"; +import { collectUploadedFileIds } from "@/lib/aidpDocumentStatus"; import { partitionAidpFiles, validateAidpFiles, } from "@/services/uploadService"; -import { getAidpUploadFailureDetails } from "@/ext_components/aidp/services/aidpUploadUtils"; -import { collectUploadedFileIds } from "@/lib/aidpDocumentStatus"; -import { useAidpGroupOptions } from "../hooks/useAidpGroupOptions"; -import { - AIDP_MODAL_STYLES, - AidpKnowledgeBaseBasicFields, - AidpKnowledgeBaseModalFooter, - AidpKnowledgeBaseModalHeader, - AidpKnowledgeBasePermissionFields, -} from "./AidpKnowledgeBaseModalParts"; +import { useGroupList } from "@/hooks/group/useGroupList"; +import { useAuthorizationContext } from "@/components/providers/AuthorizationProvider"; const { Dragger } = Upload; +// Preferred VLM model name when present in AIDP's available list. +// Falls back to the first model in the list if this specific one is absent. const PREFERRED_VLM_MODEL = "Qwen3-VL-8B-Instruct"; -type AidpPermission = "EDIT" | "READ_ONLY" | "PRIVATE"; +/** + * Default AIDP knowledge base configuration. + * Aligned with sdk/nexent/core/knowledge_base/config.py (build_create_payload defaults). + * + * Required fields per AIDP schema: + * chunk_token_num (> 0), chunk_overlap_num (>= 0) + * Reference fills the rest (is_personal, topk, similarity, smartsplit, caption_enable). + * ``vlm_model`` is no longer a hardcoded constant — it is resolved at runtime + * from the list of models AIDP advertises as applicable to the KnowledgeBase + * application (see ``useQuery(["aidp-models"])``). + */ const AIDP_CREATE_DEFAULTS = { chunk_token_num: 1024, chunk_overlap_num: 128, @@ -58,108 +73,10 @@ const AIDP_CREATE_DEFAULTS = { topk: 10, similarity: 0.0, smartsplit: 1, + // caption_enable: int 0/1. caption_enable: 0, }; -const validateAidpCreateFiles = ( - files: File[], - t: TFunction, - notify: typeof message -) => { - if (files.length === 0) return true; - const validation = validateAidpFiles(files); - if (validation.valid.length === files.length) return true; - partitionAidpFiles(files, t, notify); - return false; -}; - -const getAidpCreatePermissionValues = ( - isUser: boolean, - values: { ingroup_permission?: string; group_ids?: unknown } -) => { - const configuredPermission = values.ingroup_permission; - let permission: AidpPermission = "READ_ONLY"; - if (isUser) { - permission = "PRIVATE"; - } else if ( - configuredPermission === "EDIT" || - configuredPermission === "READ_ONLY" || - configuredPermission === "PRIVATE" - ) { - permission = configuredPermission; - } - const groupIds = - isUser || permission === "PRIVATE" - ? [] - : Array.isArray(values.group_ids) - ? values.group_ids - : []; - return { permission, groupIds }; -}; - -const showAidpCreateUploadResult = ( - result: AidpUploadResponse, - language: string, - t: TFunction -) => { - const failureDetails = getAidpUploadFailureDetails( - result.failed_list, - language, - t("aidpKnowledge.uploadFailed") - ); - const failureLines = failureDetails.map((detail, index) => ( -
{detail}
- )); - const allFailedContent = - failureLines.length > 0 ? ( - failureLines - ) : ( -
{t("aidpKnowledge.uploadFailed")}
- ); - - if (result.summary.failed > 0 && result.summary.success === 0) { - message.warning( -
-
{t("aidpKnowledge.createKbSuccess")}
- {allFailedContent} -
- ); - return; - } - if (result.summary.failed > 0) { - message.info( -
-
{t("aidpKnowledge.createKbSuccess")}
-
- {t("aidpKnowledge.uploadPartial", { - success: result.summary.success, - failed: result.summary.failed, - })} -
- {failureLines} -
- ); - return; - } - message.success( - `${t("aidpKnowledge.createKbSuccess")} | ${t( - "aidpKnowledge.uploadSuccess", - { count: result.summary.success } - )}` - ); -}; - -const getAidpCreateErrorReason = ( - error: unknown, - knowledgeBaseCreated: boolean, - t: TFunction -) => { - if (error instanceof Error && error.message.trim()) return error.message; - return knowledgeBaseCreated - ? t("aidpKnowledge.uploadFailed") - : t("aidpKnowledge.createKbFailed"); -}; - interface AidpCreateKbModalProps { open: boolean; existingKbs: AidpKnowledgeBaseItem[]; @@ -178,79 +95,140 @@ const AidpCreateKbModal: React.FC = ({ }) => { const { t, i18n } = useTranslation(); const [form] = Form.useForm(); + const [current, setCurrent] = useState(0); const [loading, setLoading] = useState(false); const [fileList, setFileList] = useState([]); const fileListRef = useRef([]); - const pendingFilesRef = useRef([]); - const rafIdRef = useRef(null); useEffect(() => { fileListRef.current = fileList; }, [fileList]); - const { isUser, canConfigureGroupPermissions, groupOptions } = - useAidpGroupOptions(); + // Antd fires beforeUpload once per file in a multi-select batch. + // The `newFiles` array may-or-may-not be the same reference across the N + // calls (behavior differs between and and antd versions), + // so we cannot rely on reference-equality for a single-call-per-batch guard. + // Instead, we collect each file in beforeUpload and schedule a single + // requestAnimationFrame flush that runs validate/add once per batch. + // This guarantees partitionAidpFiles + toast execute exactly ONCE per + // user selection, regardless of antd's internal dispatch count. + const pendingFilesRef = useRef([]); + const rafIdRef = useRef(null); + // Load the tenant's groups so the user can pick which groups may access + // the new KB. When no tenant context is available we fall back to an + // empty list and disable the group picker. + // NOTE: ``useAuthorizationContext`` is required — it is the only context + // that exposes ``user: User | null`` (with ``tenantId``). The similarly- + // named ``useAuthenticationContext`` only carries ``session`` and the + // plain ``useAuthentication`` hook doesn't carry ``user`` at all. + const { user } = useAuthorizationContext(); + const isUser = user?.role === USER_ROLES.USER; + const canConfigureGroupPermissions = !!user && !isUser; + const tenantId = user?.tenantId ?? null; + const { data: groupListData } = useGroupList( + canConfigureGroupPermissions ? tenantId : null + ); + const groupOptions = useMemo( + () => + (groupListData?.groups ?? []).map((g) => ({ + value: g.group_id, + label: g.group_name, + })), + [groupListData] + ); + const [formValues, setFormValues] = useState<{ + name: string; + description?: string; + vlm_model?: string; + chunk_token_num: number; + chunk_overlap_num: number; + caption_enable: number; + ingroup_permission: "EDIT" | "READ_ONLY" | "PRIVATE"; + group_ids: number[]; + }>({ + name: "", + chunk_token_num: AIDP_CREATE_DEFAULTS.chunk_token_num, + chunk_overlap_num: AIDP_CREATE_DEFAULTS.chunk_overlap_num, + caption_enable: AIDP_CREATE_DEFAULTS.caption_enable, + ingroup_permission: isUser ? "PRIVATE" : "READ_ONLY", + group_ids: [], + }); + + // Drive the vlm_model dropdown's visibility off the live Switch value. + // useWatch gives us a re-render whenever caption_enable toggles, without + // forcing the user to manually sync the form value to local state. const captionEnabled = Form.useWatch("caption_enable", form); + // Track in-group permission live so the group_ids picker can be disabled + // at PRIVATE without calling Form.useWatch inside a conditional sub-render + // (which would violate the Rules of Hooks). const ingroupPermission = Form.useWatch("ingroup_permission", form); useEffect(() => { if (!open) return; form.setFieldsValue({ - chunk_token_num: AIDP_CREATE_DEFAULTS.chunk_token_num, - chunk_overlap_num: AIDP_CREATE_DEFAULTS.chunk_overlap_num, - caption_enable: AIDP_CREATE_DEFAULTS.caption_enable === 1, ingroup_permission: isUser ? "PRIVATE" : "READ_ONLY", group_ids: [], }); + setFormValues((previous) => ({ + ...previous, + ingroup_permission: isUser ? "PRIVATE" : previous.ingroup_permission, + group_ids: isUser ? [] : previous.group_ids, + })); }, [form, isUser, open]); + // Fetch applicable VLM models from AIDP. Only run when modal is open to + // avoid hitting the (relatively slow) admin endpoint unnecessarily. const { data: vlmModelsData, isLoading: vlmModelsLoading } = useQuery({ queryKey: ["aidp-models", "llm", "KnowledgeBase"], queryFn: () => aidpKnowledgeService.listModels("llm", "KnowledgeBase"), enabled: open, - staleTime: 5 * 60 * 1000, + staleTime: 5 * 60 * 1000, // 5 min }); const vlmModelOptions = useMemo(() => { const models: AidpModelItem[] = vlmModelsData?.models ?? []; return models - .map((model) => model.model_name) - .filter((name): name is string => Boolean(name)); + .map((m) => m.model_name) + .filter( + (name): name is string => typeof name === "string" && name.length > 0 + ); }, [vlmModelsData]); + // Resolve the default VLM model: prefer the hardcoded + // PREFERRED_VLM_MODEL if present, otherwise the first in the list. + // Falls back to PREFERRED_VLM_MODEL (sent to AIDP as-is) when the + // models endpoint returns empty, matching the previous behavior. const defaultVlmModel = useMemo(() => { if (vlmModelOptions.length === 0) return PREFERRED_VLM_MODEL; - return vlmModelOptions.includes(PREFERRED_VLM_MODEL) - ? PREFERRED_VLM_MODEL - : vlmModelOptions[0]; + if (vlmModelOptions.includes(PREFERRED_VLM_MODEL)) + return PREFERRED_VLM_MODEL; + return vlmModelOptions[0]; }, [vlmModelOptions]); + // Pre-populate vlm_model on the form whenever the default is resolved, + // so the user sees a meaningful default on first open. useEffect(() => { - if (!open || !defaultVlmModel) return; - const currentModel = form.getFieldValue("vlm_model"); - if (!currentModel || !vlmModelOptions.includes(currentModel)) { + if (!open) return; + if (!defaultVlmModel) return; + const current = form.getFieldValue("vlm_model"); + if (!current || !vlmModelOptions.includes(current)) { form.setFieldValue("vlm_model", defaultVlmModel); } - }, [defaultVlmModel, form, open, vlmModelOptions]); + }, [open, defaultVlmModel, vlmModelOptions, form]); + // Duplicate name check against existing KBs const existingNames = useMemo( () => new Set( (existingKbs || []) .map((kb) => kb.kds_name?.toLowerCase().trim()) - .filter((name): name is string => Boolean(name)) + .filter((n): n is string => !!n) ), [existingKbs] ); - const handleSubmit = async () => { - let knowledgeBaseCreated = false; - // Files AIDP accepted while creating the KB. Reported to the parent so it - // can keep refreshing the document list until they finish processing. - let uploadedFileIds: string[] = []; - let createdKnowledgeBase: AidpKnowledgeBaseItem | null = null; - + const handleNext = async () => { try { const values = await form.validateFields(); const name = values.name.trim(); @@ -260,63 +238,173 @@ const AidpCreateKbModal: React.FC = ({ return; } - if (!validateAidpCreateFiles(fileList, t, message)) return; - - setLoading(true); - const { permission, groupIds } = getAidpCreatePermissionValues( - isUser, - values - ); - const captionEnable = values.caption_enable ? 1 : 0; - - const created = await aidpKnowledgeService.createKb({ + // Save form values before fields unmount + setFormValues({ name, - description: values.description?.trim() || "", + description: values.description?.trim() || undefined, + vlm_model: values.caption_enable + ? values.vlm_model || defaultVlmModel || undefined + : "", chunk_token_num: values.chunk_token_num ?? AIDP_CREATE_DEFAULTS.chunk_token_num, chunk_overlap_num: values.chunk_overlap_num ?? AIDP_CREATE_DEFAULTS.chunk_overlap_num, + caption_enable: values.caption_enable ? 1 : 0, + // The permission select is disabled at PRIVATE so users cannot pick + // group_ids while PRIVATE; we always coerce to [] for safety. + ingroup_permission: isUser + ? "PRIVATE" + : (values.ingroup_permission ?? "READ_ONLY"), + group_ids: + isUser || (values.ingroup_permission ?? "READ_ONLY") === "PRIVATE" + ? [] + : Array.isArray(values.group_ids) + ? values.group_ids + : [], + }); + setCurrent(1); + } catch { + // form validation error, do nothing + } + }; + + const handleBack = () => { + // Restore formValues into the Form when remounting Step 0, + // since antd Form clears field values when the Form is unmounted. + form.setFieldsValue(formValues); + setCurrent(0); + }; + + const handleSubmit = async (skipUpload: boolean) => { + let knowledgeBaseCreated = false; + // Files accepted by AIDP while creating the KB. Reported to the parent so + // it can keep refreshing the document list until they finish processing. + let uploadedFileIds: string[] = []; + let createdKdsId = ""; + let createdKnowledgeBase: AidpKnowledgeBaseItem | null = null; + try { + if (!formValues.name?.trim()) { + message.error(t("aidpKnowledge.kbNameRequired")); + setCurrent(0); + return; + } + setLoading(true); + + const permission = isUser ? "PRIVATE" : formValues.ingroup_permission; + const groupIds = isUser ? [] : formValues.group_ids; + + // Defense-in-depth: re-validate every file in case beforeUpload was bypassed + if (!skipUpload && fileList.length > 0) { + const validation = validateAidpFiles(fileList); + if (validation.valid.length !== fileList.length) { + setLoading(false); + partitionAidpFiles(fileList, t, message); + return; + } + } + + // Step 1: Create KB + // Aligned with sdk/nexent/core/knowledge_base/mapper.py#build_create_payload + const created = await aidpKnowledgeService.createKb({ + name: formValues.name.trim(), + description: formValues.description || "", + chunk_token_num: formValues.chunk_token_num, + chunk_overlap_num: formValues.chunk_overlap_num, embedding_model: AIDP_CREATE_DEFAULTS.embedding_model, - vlm_model: captionEnable - ? values.vlm_model || defaultVlmModel || "" - : "", + vlm_model: + formValues.caption_enable === 1 + ? formValues.vlm_model || defaultVlmModel || "" + : "", is_personal: AIDP_CREATE_DEFAULTS.is_personal, topk: AIDP_CREATE_DEFAULTS.topk, similarity: AIDP_CREATE_DEFAULTS.similarity, smartsplit: AIDP_CREATE_DEFAULTS.smartsplit, - caption_enable: captionEnable, + caption_enable: formValues.caption_enable, + // v7.1: forward in-group permission + groups to the backend so the + // knowledge-base permission row is created in lockstep with the KB. ingroup_permission: permission, group_ids: groupIds, }); knowledgeBaseCreated = true; + createdKdsId = String(created.kds_id || ""); createdKnowledgeBase = { ...created, - kds_id: String(created.kds_id || ""), - kds_name: created.kds_name || name, - description: created.description ?? values.description?.trim() ?? "", + kds_id: createdKdsId, + kds_name: created.kds_name || formValues.name.trim(), + description: created.description ?? formValues.description ?? "", permission: "EDIT", ingroup_permission: permission, - group_ids: groupIds, + group_ids: permission === "PRIVATE" ? [] : groupIds, resource_status: "ACTIVE", - is_multimodal: captionEnable === 1, + is_multimodal: formValues.caption_enable === 1, }; - if (fileList.length > 0 && created.kds_id) { + // Step 2: Upload files (if any and not skipped) + if (!skipUpload && fileList.length > 0 && created.kds_id) { const result = await aidpKnowledgeService.uploadDocs( created.kds_id, fileList ); - showAidpCreateUploadResult(result, i18n.language, t); uploadedFileIds = collectUploadedFileIds(result.success_list); + + const failureDetails = result.failed_list.map((item) => { + const reason = i18n.language.startsWith("zh") + ? item.reason_zh || item.reason_en + : item.reason_en || item.reason_zh; + return `${item.file_name}: ${reason || t("aidpKnowledge.uploadFailed")}`; + }); + const failureLines = failureDetails.map((detail, index) => ( +
{detail}
+ )); + + if (result.summary.failed > 0 && result.summary.success === 0) { + message.warning( +
+
{t("aidpKnowledge.createKbSuccess")}
+ {failureLines.length > 0 ? ( + failureLines + ) : ( +
{t("aidpKnowledge.uploadFailed")}
+ )} +
+ ); + } else if (result.summary.failed > 0) { + message.info( +
+
{t("aidpKnowledge.createKbSuccess")}
+
+ {t("aidpKnowledge.uploadPartial", { + success: result.summary.success, + failed: result.summary.failed, + })} +
+ {failureLines} +
+ ); + } else { + message.success( + t("aidpKnowledge.createKbSuccess") + + " | " + + t("aidpKnowledge.uploadSuccess", { + count: result.summary.success, + }) + ); + } } else { message.success(t("aidpKnowledge.createKbSuccess")); } handleReset(); - if (createdKnowledgeBase) + if (createdKnowledgeBase) { onSuccess(createdKnowledgeBase, uploadedFileIds); + } } catch (error) { - const reason = getAidpCreateErrorReason(error, knowledgeBaseCreated, t); + const reason = + error instanceof Error && error.message.trim() + ? error.message + : knowledgeBaseCreated + ? t("aidpKnowledge.uploadFailed") + : t("aidpKnowledge.createKbFailed"); message.error( knowledgeBaseCreated ? `${t("aidpKnowledge.createKbSuccess")} | ${reason}` @@ -324,8 +412,9 @@ const AidpCreateKbModal: React.FC = ({ ); if (knowledgeBaseCreated) { handleReset(); - if (createdKnowledgeBase) + if (createdKnowledgeBase) { onSuccess(createdKnowledgeBase, uploadedFileIds); + } } } finally { setLoading(false); @@ -334,12 +423,21 @@ const AidpCreateKbModal: React.FC = ({ const handleReset = () => { form.resetFields(); + setCurrent(0); setFileList([]); pendingFilesRef.current = []; if (rafIdRef.current !== null) { cancelAnimationFrame(rafIdRef.current); rafIdRef.current = null; } + setFormValues({ + name: "", + chunk_token_num: AIDP_CREATE_DEFAULTS.chunk_token_num, + chunk_overlap_num: AIDP_CREATE_DEFAULTS.chunk_overlap_num, + caption_enable: AIDP_CREATE_DEFAULTS.caption_enable, + ingroup_permission: isUser ? "PRIVATE" : "READ_ONLY", + group_ids: [], + }); }; const handleCancel = () => { @@ -347,253 +445,330 @@ const AidpCreateKbModal: React.FC = ({ onCancel(); }; - const addFiles = (files: File[]) => { - const currentFiles = fileListRef.current; - const existing = new Set(currentFiles.map((file) => file.name)); - const uniqueFiles = files.filter((file) => !existing.has(file.name)); - const { valid } = partitionAidpFiles( - uniqueFiles, - t, - message, - currentFiles.length - ); - if (valid.length > 0) setFileList([...currentFiles, ...valid]); - }; + // ---- Render steps ---- + + const renderStep0 = () => ( + <> +
+ + + + + + + + {canConfigureGroupPermissions && ( + <> + {/* USER creation is always a personal PRIVATE KB. */} + + + + + )} + + + {t("aidpKnowledge.createCaptionEnable")} + + + + + } + > + + + + {/* VLM model picker is only relevant when multimodal captioning is + enabled. Hide the dropdown entirely when the Switch is off so + users aren't shown an irrelevant choice, and so the backend + receives an empty ``vlm_model`` (see handleSubmit). */} + {captionEnabled && ( + + {t("aidpKnowledge.createVlmModel")} + + + + + } + > + ({ - label: name, - value: name, - }))} - filterOption={(input, option) => - (option?.label as string) - ?.toLowerCase() - .includes(input.toLowerCase()) ?? false - } - /> - - )} - - ), - }, - ]} - /> -
- + + + {current === 0 && ( + + )} + {current === 1 && fileList.length > 0 && ( + + )} + {current === 1 && ( + + )} + + + } + > + + + {current === 0 && renderStep0()} + {current === 1 && renderStep1()} ); }; diff --git a/frontend/ext_components/aidp/components/AidpDocumentList.tsx b/frontend/ext_components/aidp/components/AidpDocumentList.tsx index 8e8e7a037..b8131da27 100644 --- a/frontend/ext_components/aidp/components/AidpDocumentList.tsx +++ b/frontend/ext_components/aidp/components/AidpDocumentList.tsx @@ -1,9 +1,9 @@ -import React, { useCallback, useRef, useState } from "react"; +import React, { useState, useCallback, useRef } from "react"; import { useTranslation } from "react-i18next"; -import { Button, Modal, Pagination, Tag, Upload, message, Tooltip } from "antd"; +import { Button, Pagination, Tag, Upload, message, Tooltip } from "antd"; import { - FileTextOutlined, + UploadOutlined, InboxOutlined, ReloadOutlined, } from "@ant-design/icons"; @@ -12,7 +12,6 @@ import type { AidpKnowledgeBaseItem } from "@/types/agentConfig"; import type { AidpDocumentItem } from "@/ext_components/aidp/services/aidpKnowledgeService"; import aidpKnowledgeService from "@/ext_components/aidp/services/aidpKnowledgeService"; import { AIDP_ACCEPT_STRING } from "@/const/knowledgeBase"; -import log from "@/lib/logger"; import { AIDP_DOC_IN_PROGRESS_STATUSES, AIDP_DOCUMENT_STATUS, @@ -20,7 +19,6 @@ import { normalizeAidpDocStatus, } from "@/lib/aidpDocumentStatus"; import { partitionAidpFiles } from "@/services/uploadService"; -import { getAidpUploadFailureDetails } from "@/ext_components/aidp/services/aidpUploadUtils"; const { Dragger } = Upload; @@ -53,6 +51,28 @@ const isDuplicateUploadReason = ( ); }; +/** Table cell showing a document name above its AIDP file id. */ +const DocumentNameCell: React.FC<{ fileName: string; fileInoNo: string }> = ({ + fileName, + fileInoNo, +}) => ( + + {/* `max-width` on a table cell is ignored by browsers, so the clamp must + live on the inner div; the full value is surfaced on hover through an + antd Tooltip instead of a width measurement. */} + +
+ {fileName} +
+
+ +
+ {fileInoNo} +
+
+ +); + /** * Labels for the in-progress statuses, which all render as a blue tag. * @@ -83,13 +103,13 @@ const DocumentStatusCell: React.FC<{ status?: string }> = ({ status }) => { const normalized = normalizeAidpDocStatus(status); if (!normalized) { - return -; + return -; } if (AIDP_DOC_IN_PROGRESS_STATUSES.includes(normalized)) { const labelKey = IN_PROGRESS_STATUS_LABELS[normalized]; return ( - + {labelKey ? t(labelKey) : status} ); @@ -97,7 +117,7 @@ const DocumentStatusCell: React.FC<{ status?: string }> = ({ status }) => { if (normalized === AIDP_DOCUMENT_STATUS.COMPLETED) { return ( - + {t("aidpKnowledge.docStatusCompleted")} ); @@ -105,50 +125,34 @@ const DocumentStatusCell: React.FC<{ status?: string }> = ({ status }) => { if (normalized === AIDP_DOCUMENT_STATUS.FAILED) { return ( - + {t("aidpKnowledge.docStatusFailed")} ); } return ( - + {status} ); }; -const resolveDownloadFilename = (response: Response, fallback: string) => { - const contentDisposition = response.headers.get("content-disposition") || ""; - const encodedName = /filename\*=UTF-8''([^;]+)/i.exec( - contentDisposition - )?.[1]; - if (encodedName) { - try { - return decodeURIComponent(encodedName); - } catch { - // Use the regular filename or document name when decoding fails. - } - } - const plainName = /filename="?([^";]+)"?/i.exec(contentDisposition)?.[1]; - return plainName || response.headers.get("x-file-name") || fallback; -}; - interface AidpDocumentListProps { activeKb: AidpKnowledgeBaseItem | null; documents: AidpDocumentItem[]; totalDocs: number; - /** True when `totalDocs` came from the AIDP Count API. */ + /** True when `totalDocs` came from the AIDP Count API; when false the + * total is a fallback estimate and "共 N 条" should be suppressed. */ totalReliable: boolean; hasMore: boolean; isLoading: boolean; currentPage: number; pageSize: number; onPageChange: (page: number) => void; - /** Called after documents change, with the ids AIDP returned for any newly - * accepted uploads (an empty list for deletions/refreshes). The parent uses - * them to keep refreshing the list until each uploaded file reports a - * terminal processing status. */ + /** Called after an upload is accepted, with the ids AIDP returned for the + * accepted files. The parent uses them to keep refreshing the list until + * each uploaded file reports a terminal processing status. */ onDocsUploaded: (uploadedFileIds: string[]) => void; onRefresh: () => void; } @@ -168,10 +172,6 @@ const AidpDocumentList: React.FC = ({ }) => { const { t, i18n } = useTranslation(); const [uploading, setUploading] = useState(false); - const [deleting, setDeleting] = useState(false); - const [downloadingFileUuid, setDownloadingFileUuid] = useState( - null - ); // Antd fires beforeUpload once per file in a multi-select batch. // The `fileList` array may-or-may-not be the same reference across the N // calls (behavior differs between and and antd versions), @@ -181,85 +181,10 @@ const AidpDocumentList: React.FC = ({ const pendingFilesRef = useRef([]); const rafIdRef = useRef(null); - const isUnavailable = - activeKb?.resource_status === "UNAVAILABLE" || - activeKb?.resource_status === "ORPHANED"; - const canDeleteDocuments = - !!activeKb && !isUnavailable && activeKb.permission === "EDIT"; - const canDownloadDocuments = - !!activeKb && - !isUnavailable && - (activeKb.permission === "EDIT" || activeKb.permission === "READ_ONLY"); - - const handleDownload = useCallback( - async (document: AidpDocumentItem) => { - if (!activeKb || !document.file_uuid) return; - setDownloadingFileUuid(document.file_uuid); - try { - const response = await aidpKnowledgeService.downloadDoc( - activeKb.kds_id, - document.file_uuid - ); - const blob = await response.blob(); - const downloadUrl = URL.createObjectURL(blob); - const link = window.document.createElement("a"); - link.href = downloadUrl; - link.download = resolveDownloadFilename(response, document.file_name); - link.click(); - URL.revokeObjectURL(downloadUrl); - message.success(t("aidpKnowledge.downloadSuccess")); - } catch (error) { - log.error("Failed to download AIDP document:", error); - message.error(t("aidpKnowledge.downloadFailed")); - } finally { - setDownloadingFileUuid(null); - } - }, - [activeKb, t] - ); - - const handleDelete = useCallback( - (document: AidpDocumentItem) => { - if (!activeKb || !document.file_uuid) return; - Modal.confirm({ - title: t("aidpKnowledge.confirmDeleteDocTitle"), - content: t("aidpKnowledge.confirmDeleteDocContent"), - okText: t("common.confirm"), - cancelText: t("common.cancel"), - okButtonProps: { danger: true }, - centered: true, - onOk: async () => { - setDeleting(true); - try { - const result = await aidpKnowledgeService.removeDoc( - activeKb.kds_id, - document.file_uuid - ); - if (result.summary.success > 0) { - message.success(t("aidpKnowledge.deleteDocSuccess")); - } else { - message.error(t("aidpKnowledge.deleteDocFailed")); - } - if (result.summary.success > 0) { - // Deletions are not uploads: nothing to wait for, so the parent - // simply refreshes the list. - onDocsUploaded([]); - } - } catch (error) { - log.error("Failed to delete AIDP document:", error); - message.error(t("aidpKnowledge.deleteDocFailed")); - } finally { - setDeleting(false); - } - }, - }); - }, - [activeKb, onDocsUploaded, t] - ); - const handleUpload = useCallback( async (fileList: File[]) => { - if (!activeKb || fileList.length === 0) return; + if (!activeKb) return; + if (fileList.length === 0) return; setUploading(true); try { @@ -268,21 +193,20 @@ const AidpDocumentList: React.FC = ({ fileList ); - // A rejected duplicate is not an error the user can debug, so give it a - // dedicated message instead of echoing AIDP's "please rename or delete" - // instruction, which is not actionable in this dialog. Every other - // failure keeps the shared reason formatter. const failureDetails = result.failed_list.map((item) => { + // A duplicate is not an error the user can debug, so give it a + // dedicated message instead of echoing AIDP's "please rename or + // delete" instruction, which is not actionable in this dialog. if (isDuplicateUploadReason(item.reason_zh, item.reason_en)) { return t("aidpKnowledge.uploadDuplicateFile", { fileName: item.file_name, }); } - return getAidpUploadFailureDetails( - [item], - i18n.language, - t("aidpKnowledge.uploadFailed") - )[0]; + + const reason = i18n.language.startsWith("zh") + ? item.reason_zh || item.reason_en + : item.reason_en || item.reason_zh; + return `${item.file_name}: ${reason || t("aidpKnowledge.uploadFailed")}`; }); const failureLines = failureDetails.map((detail, index) => (
{detail}
@@ -328,6 +252,7 @@ const AidpDocumentList: React.FC = ({ [activeKb, i18n.language, onDocsUploaded, t] ); + // Format file size for display const formatSize = (bytes?: number): string => { if (!bytes || bytes === 0) return "-"; if (bytes < 1024) return `${bytes} B`; @@ -337,231 +262,205 @@ const AidpDocumentList: React.FC = ({ return `${(bytes / (1024 * 1024 * 1024)).toFixed(1)} GB`; }; - const effectiveTotal = totalReliable - ? totalDocs - : hasMore - ? currentPage * pageSize + 1 - : currentPage * pageSize; - - const canUpload = - !!activeKb && !isUnavailable && activeKb.permission === "EDIT"; - - const renderUploadArea = () => { - if (!canUpload) { - const reasonKey = !activeKb - ? "aidpKnowledge.uploadNoKb" - : isUnavailable - ? "aidpKnowledge.uploadKbUnavailable" - : "aidpKnowledge.uploadReadOnly"; - return ( -
-

{t(reasonKey)}

-
- ); - } - - return ( - { - pendingFilesRef.current.push(_file); - if (rafIdRef.current === null) { - rafIdRef.current = requestAnimationFrame(() => { - const batch = pendingFilesRef.current; - pendingFilesRef.current = []; - rafIdRef.current = null; - - const { valid } = partitionAidpFiles(batch, t, message); - if (valid.length > 0) void handleUpload(valid); - }); - } - return false; - }} - disabled={uploading} - className="!rounded-xl !border-blue-200 !bg-blue-50/30" - > -

- -

-

- {uploading - ? t("aidpKnowledge.uploading") - : t("aidpKnowledge.uploadHint")} -

-
-
{t("aidpKnowledge.uploadHintCount")}
-
{t("aidpKnowledge.uploadHintSize")}
-
- {t("aidpKnowledge.uploadHintFormats")} -
-
-
- ); - }; - return ( -
-
-
-
- -
-
- {/* An empty `title` leaves the Tooltip inert, so it is safe - while the knowledge base detail is still loading. */} +
+ {/* Header */} +
+
+
-

+

{activeKb?.kds_name || ""} -

+
- + {t("aidpKnowledge.tagDocs", { count: totalDocs })}
+ +
- -
-
+ {/* Document table */} +
{isLoading ? ( -
+
-
-

+

+

{t("aidpKnowledge.loadingDocs")}

) : documents.length > 0 ? ( -
+
- + - - - - - - - + {documents.map((doc) => ( - - - + + - - - ))}
+ {t("aidpKnowledge.docFileName")} + {t("aidpKnowledge.docType")} + {t("aidpKnowledge.docStatus")} + {t("aidpKnowledge.docSize")} + {t("aidpKnowledge.docCreatedAt")} - {t("aidpKnowledge.docActions")} -
- {/* Long file names are clamped by CSS on a - width-bounded inner element (`max-width` on a - table cell is ignored by browsers, so the clamp - must live on the div itself) and the full value - is surfaced through an antd Tooltip on hover, - matching the local knowledge-base list. */} - -
- {doc.file_name} -
-
- -
- {doc.file_ino_no} -
-
-
+
{doc.file_type || "-"} + {formatSize(doc.file_size)} + {doc.created_at ? new Date(doc.created_at).toLocaleString() : "-"} -
- {canDownloadDocuments && ( - - )} - {canDeleteDocuments && ( - - )} -
-
) : ( -
+
{t("aidpKnowledge.noDocuments")}
)}
- {documents.length > 0 && ( -
- t("aidpKnowledge.showTotal", { count }) - : undefined - } - size="small" - /> -
- )} - -
- {renderUploadArea()} + {/* Server-side pagination. + AIDP exposes a dedicated Count API for documents which the backend + now calls alongside the list request. When Count succeeds, + `totalReliable` is true and we display the full pagination (page + numbers + "共 N 条"). When Count fails (e.g. the endpoint is not + available on a particular AIDP instance), `totalReliable` is false + and we fall back to simple prev/next mode without a total, using + `has_more` to decide whether the next-page button should enable. */} + {documents.length > 0 && + (() => { + // When total is unreliable we still need antd to know when to + // enable "next": set total just past the current page if there is + // a next page, otherwise clamp to the current page end. + const effectiveTotal = totalReliable + ? totalDocs + : hasMore + ? currentPage * pageSize + 1 + : currentPage * pageSize; + return ( +
+ t("aidpKnowledge.showTotal", { count: total }) + : undefined + } + size="small" + /> +
+ ); + })()} + + {/* Upload area — gated by ``activeKb.permission`` and ``resource_status``. + + Per v7.1 §7.1, READ_ONLY callers may view existing documents but + must not be able to upload. UNAVAILABLE / ORPHANED KBs are + read-only regardless of permission because the AIDP backend cannot + service the request. The container is replaced with a hint instead + of disabling the Dragger so the visual structure stays consistent + and screen-reader users get an explicit reason. */} +
+ {(() => { + const isUnavailable = + activeKb?.resource_status === "UNAVAILABLE" || + activeKb?.resource_status === "ORPHANED"; + const canUpload = + !!activeKb && !isUnavailable && activeKb.permission === "EDIT"; + if (!canUpload) { + const reasonKey = !activeKb + ? "aidpKnowledge.uploadNoKb" + : isUnavailable + ? "aidpKnowledge.uploadKbUnavailable" + : "aidpKnowledge.uploadReadOnly"; + return ( +
+

{t(reasonKey)}

+
+ ); + } + return ( + { + // Queue the file and defer validation + upload until the + // synchronous batch of beforeUpload calls finishes. Each batch + // flushes in a single frame so toasts and handleUpload run once. + pendingFilesRef.current.push(_file); + if (rafIdRef.current === null) { + rafIdRef.current = requestAnimationFrame(() => { + const batch = pendingFilesRef.current; + pendingFilesRef.current = []; + rafIdRef.current = null; + + const { valid } = partitionAidpFiles(batch, t, message); + if (valid.length > 0) { + handleUpload(valid); + } + }); + } + return false; + }} + disabled={uploading} + > +

+ +

+

+ {uploading + ? t("aidpKnowledge.uploading") + : t("aidpKnowledge.uploadHint")} +

+
+
{t("aidpKnowledge.uploadHintCount")}
+
{t("aidpKnowledge.uploadHintSize")}
+
+ {t("aidpKnowledge.uploadHintFormats")} +
+
+
+ ); + })()}
); diff --git a/frontend/ext_components/aidp/components/AidpKnowledgeBaseModalParts.tsx b/frontend/ext_components/aidp/components/AidpKnowledgeBaseModalParts.tsx deleted file mode 100644 index 85b2d77a8..000000000 --- a/frontend/ext_components/aidp/components/AidpKnowledgeBaseModalParts.tsx +++ /dev/null @@ -1,162 +0,0 @@ -import React from "react"; -import { Button, Form, Input, Select } from "antd"; -import type { TFunction } from "i18next"; - -import { AIDP_KNOWLEDGE_BASE_NAME_PATTERN } from "@/const/knowledgeBase"; -import type { AidpGroupOption } from "../hooks/useAidpGroupOptions"; - -export const AIDP_MODAL_STYLES = { - container: { overflow: "hidden", borderRadius: 16, padding: 0 }, - body: { padding: 0 }, - footer: { - margin: 0, - padding: "12px 20px 16px", - borderTop: "1px solid #f0f0f0", - }, -}; - -interface AidpModalHeaderProps { - title: string; - subtitle: string; -} - -export const AidpKnowledgeBaseModalHeader: React.FC = ({ - title, - subtitle, -}) => ( -
-

- {title} -

-

{subtitle}

-
-); - -interface AidpModalFooterProps { - onCancel: () => void; - onSubmit: () => void; - loading: boolean; - cancelText: string; - submitText: string; -} - -export const AidpKnowledgeBaseModalFooter: React.FC = ({ - onCancel, - onSubmit, - loading, - cancelText, - submitText, -}) => ( -
- - -
-); - -interface AidpKnowledgeBaseBasicFieldsProps { - t: TFunction; -} - -export const AidpKnowledgeBaseBasicFields: React.FC< - AidpKnowledgeBaseBasicFieldsProps -> = ({ t }) => ( - <> - - - - - - - -); - -interface AidpPermissionFieldsProps { - t: TFunction; - groupOptions: AidpGroupOption[]; - ingroupPermission?: string; - showSearch?: boolean; -} - -const hasRequiredGroupIds = (permission: string, value: unknown) => - permission === "PRIVATE" || (Array.isArray(value) && value.length > 0); - -export const AidpKnowledgeBasePermissionFields: React.FC< - AidpPermissionFieldsProps -> = ({ t, groupOptions, ingroupPermission, showSearch = false }) => ( - <> - - - - -); diff --git a/frontend/ext_components/aidp/components/AidpKnowledgeConfiguration.tsx b/frontend/ext_components/aidp/components/AidpKnowledgeConfiguration.tsx index d9463c498..bef0cf94e 100644 --- a/frontend/ext_components/aidp/components/AidpKnowledgeConfiguration.tsx +++ b/frontend/ext_components/aidp/components/AidpKnowledgeConfiguration.tsx @@ -9,6 +9,7 @@ import { InfoCircleFilled } from "@ant-design/icons"; import { SETUP_PAGE_CONTAINER, TWO_COLUMN_LAYOUT, + STANDARD_CARD, } from "@/const/layoutConstants"; import { KB_SEARCH_DEBOUNCE_MS } from "@/const/knowledgeBase"; import { @@ -45,10 +46,8 @@ const AidpKnowledgeConfiguration: React.FC = () => { // which may not contain the currently active KB. `selectedKb` is the item // itself — set on selection, kept stable across list refetches. const [activeKbId, setActiveKbId] = useState(null); - const [selectedKb, setSelectedKb] = useState( - null - ); - const [, setActiveKbDetail] = useState(null); + const [selectedKb, setSelectedKb] = useState(null); + const [activeKbDetail, setActiveKbDetail] = useState(null); const [documents, setDocuments] = useState([]); const [totalDocs, setTotalDocs] = useState(0); const [docHasMore, setDocHasMore] = useState(false); @@ -67,9 +66,7 @@ const AidpKnowledgeConfiguration: React.FC = () => { // ---- Modal state ---- const [createModalOpen, setCreateModalOpen] = useState(false); const [updateModalOpen, setUpdateModalOpen] = useState(false); - const [editingKb, setEditingKb] = useState( - null - ); + const [editingKb, setEditingKb] = useState(null); // ---- Keyword search state ---- // `kbKeyword` is the raw input value and keeps the text field responsive; @@ -93,7 +90,7 @@ const AidpKnowledgeConfiguration: React.FC = () => { const result = await aidpKnowledgeService.listKbs( page, KB_PAGE_SIZE, - keyword + keyword, ); setKbs(result.value); setKbTotal(result.total_count ?? result.value.length); @@ -151,7 +148,7 @@ const AidpKnowledgeConfiguration: React.FC = () => { const result = await aidpKnowledgeService.listDocs( kbId, page, - DOC_PAGE_SIZE + DOC_PAGE_SIZE, ); const count = result.total_count ?? result.value.length; setDocuments(result.value); @@ -296,7 +293,7 @@ const AidpKnowledgeConfiguration: React.FC = () => { // Refresh list, keeping the active search filter applied fetchKbs(kbPage, debouncedKbKeyword); - } catch { + } catch (error) { appMessage.error(t("aidpKnowledge.deleteKbFailed")); } }, @@ -424,17 +421,18 @@ const AidpKnowledgeConfiguration: React.FC = () => { return (
-
- + {/* Two-column layout — content-sized cards with a single + scroll container; no card stretches to viewport height. */} +
+ {/* Left column: KB list */} { {/* Right column: Document list or empty state */} { onRefresh={handleRefreshDocs} /> ) : ( -
-
-
- +
+
+
+
+ +
+

+ {t("aidpKnowledge.selectKbTitle")} +

+

+ {t("aidpKnowledge.selectKbHint")} +

-

- {t("aidpKnowledge.selectKbTitle")} -

-

- {t("aidpKnowledge.selectKbHint")} -

)} diff --git a/frontend/ext_components/aidp/components/AidpKnowledgeList.tsx b/frontend/ext_components/aidp/components/AidpKnowledgeList.tsx index b21d0bd9f..1b5470626 100644 --- a/frontend/ext_components/aidp/components/AidpKnowledgeList.tsx +++ b/frontend/ext_components/aidp/components/AidpKnowledgeList.tsx @@ -1,19 +1,13 @@ import React, { useMemo } from "react"; import { useTranslation } from "react-i18next"; -import { Button, Input, Pagination, Tooltip } from "antd"; +import { Button, Input, Pagination, Tag, Tooltip } from "antd"; import { - BookOpen, - CircleOff, - Eye, - FolderOpen, - Glasses, - PencilRuler, - Search, - SquarePen, - Trash2, -} from "lucide-react"; -import { PlusOutlined, ReloadOutlined } from "@ant-design/icons"; + PlusOutlined, + ReloadOutlined, + SearchOutlined, +} from "@ant-design/icons"; +import { SquarePen, Trash2 } from "lucide-react"; import type { AidpKnowledgeBaseItem } from "@/types/agentConfig"; import { useGroupList } from "@/hooks/group/useGroupList"; @@ -25,16 +19,14 @@ interface AidpKnowledgeListProps { activeKbId: string | null; isLoading: boolean; total: number; - /** True when `total` came from AIDP Count API. */ + /** True when `total` came from AIDP Count API (reliable). When false we + * show a simple prev/next pagination without "共 N 条". */ totalReliable: boolean; hasMore: boolean; currentPage: number; pageSize: number; /** Raw search box value. The parent debounces it before querying, so this is - * intentionally the un-debounced text the user is currently typing. The - * filtering itself happens server-side: AIDP narrows the catalog by this - * keyword, so the page below renders exactly what the API returned and the - * reported total always matches the rendered cards. */ + * intentionally the un-debounced text the user is currently typing. */ keyword: string; onKeywordChange: (value: string) => void; onPageChange: (page: number) => void; @@ -45,20 +37,6 @@ interface AidpKnowledgeListProps { onDelete: (kb: AidpKnowledgeBaseItem) => void; } -const permissionIcon = (permission?: string) => { - const props = { size: 13, className: "text-gray-500" }; - switch (permission) { - case "EDIT": - return ; - case "READ_ONLY": - return ; - case "PRIVATE": - return ; - default: - return ; - } -}; - const AidpKnowledgeList: React.FC = ({ kbs, activeKbId, @@ -79,322 +57,235 @@ const AidpKnowledgeList: React.FC = ({ }) => { const { t } = useTranslation(); + // Load groups for the current tenant so we can render group_ids as names. + // ``useAuthorizationContext`` is the right hook here (the similarly-named + // ``useAuthenticationContext`` carries only ``session`` — no ``user`` object). const { user } = useAuthorizationContext(); const tenantId = user?.tenantId ?? null; const { data: groupListData } = useGroupList(tenantId); const groupById = useMemo(() => { const map = new Map(); - (groupListData?.groups ?? []).forEach((group) => { - map.set(group.group_id, group.group_name); + (groupListData?.groups ?? []).forEach((g) => { + map.set(g.group_id, g.group_name); }); return map; }, [groupListData]); - // The keyword filter lives on the server: the parent forwards it to AIDP, and - // filtering the returned page again here would hide results whenever AIDP's - // matching is not a plain substring of the name/description. Only the order is - // local, so the most recently updated KB stays on top of the page. - const displayedKbs = useMemo( - () => - [...kbs].sort((a, b) => { - const aTime = Date.parse(a.updated_at || a.created_at || "") || 0; - const bTime = Date.parse(b.updated_at || b.created_at || "") || 0; - return bTime - aTime; - }), - [kbs] - ); - - const getGroupNames = (groupIds?: number[]) => - (groupIds ?? []) + // Convert group ids to names, skipping any ids that don't resolve + // (e.g. the group was deleted or the list is not yet loaded). Aligned + // with the local-knowledge-base list renderer. + const getGroupNames = (groupIds: number[] | undefined): string[] => { + if (!Array.isArray(groupIds) || groupIds.length === 0) return []; + return groupIds .map((id) => groupById.get(id)) - .filter((name): name is string => Boolean(name)); - - const permissionLabel = (permission?: string) => - t(`knowledgeBase.ingroup.permission.${permission || "DEFAULT"}`); - - const renderStatusTag = (kb: AidpKnowledgeBaseItem) => { - const isUnavailable = - kb.resource_status === "UNAVAILABLE" || kb.resource_status === "ORPHANED"; - if (isUnavailable) { - return ( - - {t("aidpKnowledge.kbUnavailable")} - - ); - } - if (kb.permission === "READ_ONLY") { - return ( - - {t("aidpKnowledge.kbReadOnly")} - - ); - } - return null; + .filter((name): name is string => typeof name === "string" && name.length > 0); }; - const renderKnowledgeCard = (kb: AidpKnowledgeBaseItem) => { - const isActive = activeKbId === kb.kds_id; - const isUnavailable = - kb.resource_status === "UNAVAILABLE" || kb.resource_status === "ORPHANED"; - const canModify = kb.permission === "EDIT" && !isUnavailable; - const groupNames = getGroupNames(kb.group_ids); - const documentCount = kb.document_count ?? 0; - const chunkCount = kb.chunk_count ?? 0; - const permission = kb.ingroup_permission || "PRIVATE"; - - return ( -
onSelect(kb)} - onKeyDown={(event) => { - if (event.target !== event.currentTarget) return; - if (event.key === "Enter" || event.key === " ") { - event.preventDefault(); - onSelect(kb); - } - }} - > -
-
-
- -
-
-

- {kb.kds_name} -

- - AIDP - -
-
- -
- {canModify && ( - -
-
- -

- {kb.description?.trim() || t("knowledgeBase.description.default")} -

- -
- - {t("knowledgeBase.tag.documents", { count: documentCount })} - - - {t("knowledgeBase.tag.chunks", { count: chunkCount })} - - {kb.embedding_model && kb.embedding_model !== "default" && ( - - {kb.embedding_model} - - )} - {kb.is_multimodal && ( - - multimodal - - )} - {renderStatusTag(kb)} - - {permission === "PRIVATE" ? ( - - {permissionIcon(permission)} - {permissionLabel(permission)} - - ) : ( - groupNames.slice(0, 2).map((groupName) => ( - - {groupName} - - )) - )} - -
- -
- - {t("knowledgeBase.tag.updatedAt", { - date: - kb.updated_at || kb.created_at - ? new Date( - kb.updated_at || kb.created_at || "" - ).toLocaleDateString() - : t("aidpKnowledge.createdAtUnknown"), - })} - - - {permissionLabel(permission)} - -
-
+ // Sort alphabetically by name + const displayedKbs = useMemo(() => { + return [...kbs].sort((a, b) => + (a.kds_name || "").localeCompare(b.kds_name || "") ); - }; - - let effectiveTotal = currentPage * pageSize; - if (totalReliable) { - effectiveTotal = total; - } else if (hasMore) { - effectiveTotal += 1; - } + }, [kbs]); return ( -
-
-
-
-
- -
-
-

- {t("knowledgeBase.page.title")} -

-

- {t("knowledgeBase.page.description")} -

-
-
- -
+
+ {/* Header */} +
+
+

+ {t("aidpKnowledge.kbListTitle")} +

+
- -
-

- {t("knowledgeBase.page.all")} - - {t("knowledgeBase.page.count", { count: total })} - -

- } - value={keyword} - onChange={(event) => onKeywordChange(event.target.value)} - className="h-10 min-w-[240px] max-w-[420px] flex-1 !rounded-lg" - allowClear - /> -
+ {/* Search box. `keyword` is the raw input value and the parent debounces + it, so typing stays responsive and only the settled text triggers a + request. Clearing the field restores the unfiltered list. */} + } + value={keyword} + onChange={(e) => onKeywordChange(e.target.value)} + size="small" + />
-
- {isLoading && kbs.length === 0 ? ( -
- Loading... + {/* List */} +
+ {displayedKbs.length > 0 ? ( +
+ {displayedKbs.map((kb) => { + const isActive = activeKbId === kb.kds_id; + const isUnavailable = + kb.resource_status === "UNAVAILABLE" || + kb.resource_status === "ORPHANED"; + // Only EDIT-level callers may modify the KB or its files. + const canModify = kb.permission === "EDIT" && !isUnavailable; + + return ( +
onSelect(kb)} + > +
+
+
+

+ {kb.kds_name} +

+ {isUnavailable && ( + + {t("aidpKnowledge.kbUnavailable")} + + )} + {kb.permission === "READ_ONLY" && !isUnavailable && ( + + {t("aidpKnowledge.kbReadOnly")} + + )} +
+ {kb.description && ( +

+ {kb.description} +

+ )} +
+ {kb.ingroup_permission === "PRIVATE" && ( + + {t("knowledgeBase.ingroup.permission.PRIVATE")} + + )} + {/* Authorized user-group tags. Aligned with the local + knowledge base list: only render group names when + ``ingroup_permission !== "PRIVATE"``, each group + gets its own blue tag, and when there are no + groups to show we render nothing (no "not + authorized" fallback). Gated by the ``group:read`` + permission so users without group visibility see + the KB card cleanly without the tag area. */} + + {kb.ingroup_permission !== "PRIVATE" && + getGroupNames(kb.group_ids).map((groupName, idx) => ( + + {groupName} + + ))} + + {kb.created_at ? ( + + {t("aidpKnowledge.createdAt", { + date: new Date(kb.created_at).toLocaleDateString(), + })} + + ) : ( + {t("aidpKnowledge.createdAtUnknown")} + )} +
+
+
+ {canModify && ( + +
+
+
+ ); + })}
) : ( -
- - {displayedKbs.map(renderKnowledgeCard)} -
- )} - - {!isLoading && displayedKbs.length === 0 && ( -
+
{keyword.trim() - ? t("knowledgeBase.list.noResults") + ? t("aidpKnowledge.searchEmpty") : t("aidpKnowledge.listEmpty")}
)}
- {kbs.length > 0 && ( -
- t("aidpKnowledge.showTotal", { count }) - : undefined - } - size="small" - /> -
- )} + {/* Server-side pagination. + AIDP exposes a dedicated Count API for KBs which the backend calls + alongside the list request. When Count succeeds, `totalReliable` + is true and we display the full pagination (page numbers + + "共 N 条"). When Count fails (e.g. endpoint unavailable), we fall + back to simple prev/next mode using `has_more`. */} + {kbs.length > 0 && (() => { + const effectiveTotal = totalReliable + ? total + : (hasMore + ? currentPage * pageSize + 1 + : currentPage * pageSize); + return ( +
+ t("aidpKnowledge.showTotal", { count: total }) + : undefined + } + size="small" + /> +
+ ); + })()}
); }; diff --git a/frontend/ext_components/aidp/components/AidpUpdateKbModal.tsx b/frontend/ext_components/aidp/components/AidpUpdateKbModal.tsx index 415290f06..cb9a0317e 100644 --- a/frontend/ext_components/aidp/components/AidpUpdateKbModal.tsx +++ b/frontend/ext_components/aidp/components/AidpUpdateKbModal.tsx @@ -1,21 +1,16 @@ "use client"; -import React, { useEffect } from "react"; +import React, { useEffect, useMemo } from "react"; import { useTranslation } from "react-i18next"; -import { Collapse, Modal, Form, message } from "antd"; -import { SettingOutlined } from "@ant-design/icons"; +import { Modal, Form, Input, Select, message } from "antd"; import type { AidpKnowledgeBaseItem } from "@/types/agentConfig"; import aidpKnowledgeService from "@/ext_components/aidp/services/aidpKnowledgeService"; -import { useAidpGroupOptions } from "../hooks/useAidpGroupOptions"; -import { - AIDP_MODAL_STYLES, - AidpKnowledgeBaseBasicFields, - AidpKnowledgeBaseModalFooter, - AidpKnowledgeBaseModalHeader, - AidpKnowledgeBasePermissionFields, -} from "./AidpKnowledgeBaseModalParts"; +import { AIDP_KNOWLEDGE_BASE_NAME_PATTERN } from "@/const/knowledgeBase"; +import { useGroupList } from "@/hooks/group/useGroupList"; +import { useAuthorizationContext } from "@/components/providers/AuthorizationProvider"; +import { USER_ROLES } from "@/const/auth"; interface AidpUpdateKbModalProps { open: boolean; @@ -34,9 +29,24 @@ const AidpUpdateKbModal: React.FC = ({ const [form] = Form.useForm(); const [loading, setLoading] = React.useState(false); - const { isUser, canConfigureGroupPermissions, groupOptions } = - useAidpGroupOptions(); - const [advancedOpen, setAdvancedOpen] = React.useState(false); + // Mirror the create-modal wiring: the authorization context exposes + // ``user.tenantId``, which we feed into ``useGroupList`` to enumerate + // the tenant's groups for the access-group picker below. + const { user } = useAuthorizationContext(); + const isUser = user?.role === USER_ROLES.USER; + const canConfigureGroupPermissions = !!user && !isUser; + const tenantId = user?.tenantId ?? null; + const { data: groupListData } = useGroupList( + canConfigureGroupPermissions ? tenantId : null + ); + const groupOptions = useMemo( + () => + (groupListData?.groups ?? []).map((g) => ({ + value: g.group_id, + label: g.group_name, + })), + [groupListData] + ); const ingroupPermission = Form.useWatch("ingroup_permission", form); @@ -44,21 +54,20 @@ const AidpUpdateKbModal: React.FC = ({ // that predate the column — normalize to an empty array so the Select // (mode="multiple") receives a value shape it accepts. useEffect(() => { - if (!open) return; - setAdvancedOpen(false); - if (!knowledgeBase) return; - form.setFieldsValue({ - name: knowledgeBase.kds_name, - description: knowledgeBase.description || "", - ingroup_permission: isUser - ? "PRIVATE" - : knowledgeBase.ingroup_permission || "READ_ONLY", - group_ids: isUser - ? [] - : Array.isArray(knowledgeBase.group_ids) - ? knowledgeBase.group_ids - : [], - }); + if (open && knowledgeBase) { + form.setFieldsValue({ + name: knowledgeBase.kds_name, + description: knowledgeBase.description || "", + ingroup_permission: isUser + ? "PRIVATE" + : knowledgeBase.ingroup_permission || "READ_ONLY", + group_ids: isUser + ? [] + : Array.isArray(knowledgeBase.group_ids) + ? knowledgeBase.group_ids + : [], + }); + } }, [open, knowledgeBase, form, isUser]); const handleOk = async () => { @@ -155,74 +164,96 @@ const AidpUpdateKbModal: React.FC = ({ return ( - } > -
- -
+ - - {canConfigureGroupPermissions && ( - - setAdvancedOpen( - Array.isArray(keys) - ? keys.includes("advanced") - : keys === "advanced" - ) - } - items={[ + + + + + + {canConfigureGroupPermissions && ( + <> + - - {t("aidpKnowledge.createAdvancedOptions")} - - ), - children: ( -
- -
- ), + required: true, + message: t("aidpKnowledge.createIngroupPermissionRequired"), }, ]} - /> - )} - -
+ > + + + + )} +
); }; diff --git a/frontend/ext_components/aidp/hooks/useAidpGroupOptions.ts b/frontend/ext_components/aidp/hooks/useAidpGroupOptions.ts deleted file mode 100644 index ff0a42e36..000000000 --- a/frontend/ext_components/aidp/hooks/useAidpGroupOptions.ts +++ /dev/null @@ -1,30 +0,0 @@ -import { useMemo } from "react"; - -import { USER_ROLES } from "@/const/auth"; -import { useAuthorizationContext } from "@/components/providers/AuthorizationProvider"; -import { useGroupList } from "@/hooks/group/useGroupList"; - -export interface AidpGroupOption { - value: number; - label: string; -} - -export const useAidpGroupOptions = () => { - const { user } = useAuthorizationContext(); - const isUser = user?.role === USER_ROLES.USER; - const canConfigureGroupPermissions = Boolean(user) && !isUser; - const tenantId = user?.tenantId ?? null; - const { data: groupListData } = useGroupList( - canConfigureGroupPermissions ? tenantId : null - ); - const groupOptions = useMemo( - () => - (groupListData?.groups ?? []).map((group) => ({ - value: group.group_id, - label: group.group_name, - })), - [groupListData] - ); - - return { isUser, canConfigureGroupPermissions, groupOptions }; -}; diff --git a/frontend/ext_components/aidp/services/aidpKnowledgeService.ts b/frontend/ext_components/aidp/services/aidpKnowledgeService.ts index 80fb4830e..21b40240d 100644 --- a/frontend/ext_components/aidp/services/aidpKnowledgeService.ts +++ b/frontend/ext_components/aidp/services/aidpKnowledgeService.ts @@ -1,7 +1,7 @@ /** * AIDP Knowledge Base Management Service * - * Wraps the AIDP management backend endpoints. + * Wraps the 8 AIDP management backend endpoints. * Credentials (server_url, api_key) are read by the backend from environment variables. */ @@ -29,12 +29,15 @@ export interface AidpKbDetail { ingroup_permission?: "EDIT" | "READ_ONLY" | "PRIVATE"; group_ids?: number[]; resource_status?: - "ACTIVE" | "CREATING" | "DELETE_PENDING" | "ORPHANED" | "UNAVAILABLE"; + | "ACTIVE" + | "CREATING" + | "DELETE_PENDING" + | "ORPHANED" + | "UNAVAILABLE"; } export interface AidpDocumentItem { - file_uuid: string; - file_ino_no: number; + file_ino_no: string; file_name: string; file_size?: number; file_type?: string; @@ -67,7 +70,6 @@ export interface AidpDocumentListResponse { } export interface AidpUploadSuccessItem { - file_uuid: string; file_name: string; file_type: string; file_size: number; @@ -91,62 +93,6 @@ export interface AidpUploadResponse { failed_list: AidpUploadFailedItem[]; } -export interface AidpDocumentOperationItem { - file_uuid: string; -} - -export interface AidpDocumentRemoveResponse { - summary: { - total: number; - success: number; - failed: number; - }; - success_list: AidpDocumentOperationItem[]; - failed_list: AidpDocumentOperationItem[]; -} - -type AidpOperationSummary = { - total: number; - success: number; - failed: number; -}; - -type AidpOperationResponse = { - summary: AidpOperationSummary; - success_list: TSuccess[]; - failed_list: TFailure[]; -}; - -const normalizeAidpOperationResponse = ( - result: Partial> -): AidpOperationResponse => { - const successList: TSuccess[] = Array.isArray(result.success_list) - ? result.success_list - : []; - const failedList: TFailure[] = Array.isArray(result.failed_list) - ? result.failed_list - : []; - - return { - summary: { - total: - typeof result.summary?.total === "number" - ? result.summary.total - : successList.length + failedList.length, - success: - typeof result.summary?.success === "number" - ? result.summary.success - : successList.length, - failed: - typeof result.summary?.failed === "number" - ? result.summary.failed - : failedList.length, - }, - success_list: successList, - failed_list: failedList, - }; -}; - export interface AidpModelItem { /** Display / identifier used for the model (sent to AIDP as ``vlm_model``). */ model_name: string; @@ -421,10 +367,31 @@ class AidpKnowledgeService { } const result = (await response.json()) as Partial; - return normalizeAidpOperationResponse< - AidpUploadSuccessItem, - AidpUploadFailedItem - >(result); + const successList = Array.isArray(result.success_list) + ? result.success_list + : []; + const failedList = Array.isArray(result.failed_list) + ? result.failed_list + : []; + + return { + summary: { + total: + typeof result.summary?.total === "number" + ? result.summary.total + : successList.length + failedList.length, + success: + typeof result.summary?.success === "number" + ? result.summary.success + : successList.length, + failed: + typeof result.summary?.failed === "number" + ? result.summary.failed + : failedList.length, + }, + success_list: successList, + failed_list: failedList, + }; } /** @@ -511,46 +478,6 @@ class AidpKnowledgeService { : undefined, }; } - - /** - * Remove one document from an AIDP knowledge base. - * The AIDP API accepts an array, so the single-document UI sends one item. - */ - async removeDoc( - id: string, - fileUuid: string - ): Promise { - const url = buildUrl(API_ENDPOINTS.aidpMgmt.removeKbDocuments(id), {}); - const response = await fetchWithErrorHandling(url, { - method: "POST", - headers: { - ...getAuthHeaders(), - "Content-Type": "application/json", - }, - body: JSON.stringify({ file_uuids: [fileUuid] }), - }); - const result = - (await response.json()) as Partial; - return normalizeAidpOperationResponse< - AidpDocumentOperationItem, - AidpDocumentOperationItem - >(result); - } - - /** - * Download one document through the AIDP management backend. - */ - async downloadDoc(id: string, fileUuid: string): Promise { - const url = buildUrl(API_ENDPOINTS.aidpMgmt.downloadKbDocument(id), {}); - return fetchWithErrorHandling(url, { - method: "POST", - headers: { - ...getAuthHeaders(), - "Content-Type": "application/json", - }, - body: JSON.stringify({ file_uuid: fileUuid }), - }); - } } const aidpKnowledgeService = new AidpKnowledgeService(); diff --git a/frontend/ext_components/aidp/services/aidpUploadUtils.ts b/frontend/ext_components/aidp/services/aidpUploadUtils.ts deleted file mode 100644 index 98de81dfb..000000000 --- a/frontend/ext_components/aidp/services/aidpUploadUtils.ts +++ /dev/null @@ -1,15 +0,0 @@ -import type { AidpUploadFailedItem } from "./aidpKnowledgeService"; - -export const getAidpUploadFailureDetails = ( - failedList: AidpUploadFailedItem[], - language: string, - fallbackMessage: string -) => { - const isChinese = language.startsWith("zh"); - return failedList.map((item) => { - const reason = isChinese - ? item.reason_zh || item.reason_en - : item.reason_en || item.reason_zh; - return `${item.file_name}: ${reason || fallbackMessage}`; - }); -}; diff --git a/frontend/features/humanInteraction/client.ts b/frontend/features/humanInteraction/client.ts index 261c08f92..29547749e 100644 --- a/frontend/features/humanInteraction/client.ts +++ b/frontend/features/humanInteraction/client.ts @@ -37,6 +37,11 @@ export type HumanDecision = "answer" | "approve" | "reject" | "steer"; const base = `${API_BASE_URL}/agent/human-interactions`; +// A hung request (half-open proxy/socket) would otherwise keep the +// controller's in-flight guard stuck forever and silently drop every +// subsequent snapshot, including forced ones. Abort and surface an error. +const REQUEST_TIMEOUT_MS = 15000; + export class HumanInteractionHttpError extends Error { constructor( message: string, @@ -50,6 +55,7 @@ async function request(path: string, body?: unknown): Promise { const response = await fetchWithAuth(`${base}${path}`, { method: body === undefined ? "GET" : "POST", headers: { ...getAuthHeaders(), "Content-Type": "application/json" }, + signal: AbortSignal.timeout(REQUEST_TIMEOUT_MS), ...(body === undefined ? {} : { body: JSON.stringify(body) }), }); if (!response.ok) { diff --git a/frontend/features/humanInteraction/useHumanInteractionController.ts b/frontend/features/humanInteraction/useHumanInteractionController.ts index 9c45849c4..bd27751a6 100644 --- a/frontend/features/humanInteraction/useHumanInteractionController.ts +++ b/frontend/features/humanInteraction/useHumanInteractionController.ts @@ -19,6 +19,8 @@ const ACTIVE_STATUSES = new Set([ "RECOVERY_REQUIRED", ]); const STREAM_RECONNECT_STATUSES = new Set(["READY", "RUNNING"]); +// Safety-net snapshot poll while a run is active (see the polling effect below). +const ACTIVE_SNAPSHOT_POLL_INTERVAL_MS = 5000; const TERMINAL_STATUSES = new Set([ "COMPLETED", "FAILED", @@ -42,7 +44,7 @@ export interface HumanInteractionController { decision: HumanDecision, response: string | HumanClarificationAnswer[] ) => Promise; - refresh: () => Promise; + refresh: (force?: boolean) => Promise; } export function useHumanInteractionController({ @@ -71,6 +73,12 @@ export function useHumanInteractionController({ const refreshInFlight = useRef(false); const lastSnapshotAt = useRef(0); const MIN_SNAPSHOT_INTERVAL_MS = 3000; + // A forced refresh that arrived while another snapshot was in flight is + // re-run afterwards, so critical HITL transitions are never dropped. + const pendingForceRefresh = useRef(false); + const refreshRef = useRef< + (force?: boolean) => Promise | undefined + >(undefined); // Mirror of the latest `run` state; keeps `refresh` dependencies stable. const runRef = useRef(null); activeConversation.current = conversationId; @@ -98,7 +106,7 @@ export function useHumanInteractionController({ }; }, []); - const refresh = useCallback(async () => { + const refresh = useCallback(async (force = false) => { const requestedConversation = conversationId; if (!conversationId) { runRef.current = null; @@ -106,13 +114,16 @@ export function useHumanInteractionController({ setActiveRunId(null); return null; } - // Guard 1: coalesce while a snapshot is in flight. + // Guard 1: coalesce while a snapshot is in flight. Forced callers are + // re-run once the in-flight snapshot completes instead of being dropped. if (refreshInFlight.current) { + if (force) pendingForceRefresh.current = true; return runRef.current; } - // Guard 2: respect the minimum snapshot interval. + // Guard 2: respect the minimum snapshot interval; forced calls bypass it + // because they carry run-state transitions (e.g. ask_user suspends). const now = Date.now(); - if (now - lastSnapshotAt.current < MIN_SNAPSHOT_INTERVAL_MS) { + if (!force && now - lastSnapshotAt.current < MIN_SNAPSHOT_INTERVAL_MS) { return runRef.current; } refreshInFlight.current = true; @@ -146,9 +157,19 @@ export function useHumanInteractionController({ return null; } finally { refreshInFlight.current = false; + if (pendingForceRefresh.current) { + pendingForceRefresh.current = false; + // Routed through the latest refresh closure, which re-validates the + // active conversation and sequence before applying a snapshot. + void refreshRef.current?.(true); + } } }, [conversationId]); + useEffect(() => { + refreshRef.current = refresh; + }, [refresh]); + /** * Discovery: one-shot snapshots at conversation change and agent pause. * After a snapshot reveals a run, the SSE subscription below keeps it live. @@ -197,15 +218,29 @@ export function useHumanInteractionController({ // auto-reconnecting against a run that already finished. try { const parsed = JSON.parse(ev.data); + if (parsed.type === "human_run") { + if ( + parsed.content && + typeof parsed.content === "object" && + TERMINAL_STATUSES.has(parsed.content.status) + ) { + stoppedByUs = true; + es.close(); + void refresh(true); + return; + } + // Status transitions (e.g. ask_user → WAITING_HUMAN) carry the + // request form; they must not be dropped by the snapshot throttle — + // the event stream goes quiet afterwards, so nothing re-triggers. + void refresh(true); + return; + } if ( - parsed.type === "human_run" && - parsed.content && - typeof parsed.content === "object" && - TERMINAL_STATUSES.has(parsed.content.status) + ["human_interaction", "human_decision", "human_execution"].includes( + parsed.type + ) ) { - stoppedByUs = true; - es.close(); - void refresh(); + void refresh(true); return; } } catch { @@ -267,6 +302,20 @@ export function useHumanInteractionController({ const scopedRun = run?.conversation_id === conversationId ? run : null; const active = Boolean(scopedRun && ACTIVE_STATUSES.has(scopedRun.status)); + // Safety-net polling: the SSE stream is the primary delivery channel, but a + // silently stalled connection (hung dev proxy, half-open socket that never + // raises an error event) must not hide a pending form until a manual page + // refresh. The read-only snapshot is cheap; poll it at a low frequency while + // the scoped run is still active. Calls share the refresh() throttle and + // in-flight guard, so live SSE delivery simply coalesces with these ticks. + useEffect(() => { + if (!available || !conversationId || !active) return; + const timer = setInterval(() => { + void refresh(); + }, ACTIVE_SNAPSHOT_POLL_INTERVAL_MS); + return () => clearInterval(timer); + }, [available, conversationId, active, refresh]); + const control = useCallback( async (action: "pause" | "terminate") => { if (!run) return; diff --git a/frontend/public/locales/en/common.json b/frontend/public/locales/en/common.json index 81681aaa6..3b0cbd0bb 100644 --- a/frontend/public/locales/en/common.json +++ b/frontend/public/locales/en/common.json @@ -844,37 +844,7 @@ "knowledgeBase.error.syncFailed": "Failed to sync DataMate knowledge bases", "knowledgeBase.message.testingConnection": "Testing connection...", "knowledgeBase.message.testingSync": "Syncing knowledge bases...", - "knowledgeBase.page.title": "Knowledge Base", - "knowledgeBase.page.description": "Create, manage, and maintain your team's knowledge assets", - "knowledgeBase.page.all": "All Knowledge Bases", - "knowledgeBase.page.count": "{{count}} total", - "knowledgeBase.page.back": "Back to knowledge bases", - "knowledgeBase.create.subtitle": "Configure basic information and adjust it later at any time", - "knowledgeBase.create.field.name": "Knowledge base name", - "knowledgeBase.create.field.description": "Description", - "knowledgeBase.create.optional": "(Optional)", - "knowledgeBase.create.namePlaceholder": "e.g. Product Knowledge Hub", - "knowledgeBase.create.descriptionPlaceholder": "Briefly describe what this knowledge base contains", - "knowledgeBase.create.advancedSettings": "Advanced settings", - "knowledgeBase.create.submit": "Create and enter", - "knowledgeBase.create.field.embeddingModel": "Embedding model", - "knowledgeBase.create.field.groups": "User groups", - "knowledgeBase.create.field.permission": "Group permission", - "knowledgeBase.create.field.preserve": "Document copy", - "knowledgeBase.create.field.quota": "Storage quota", - "knowledgeBase.create.uploadTitle": "Upload documents to build your knowledge base", "knowledgeBase.list.title": "Knowledge Base List", - "knowledgeBase.personalCapacity.title": "Personal knowledge base capacity", - "knowledgeBase.personalCapacity.withQuota": "Used {{used}} / Total {{quota}}", - "knowledgeBase.personalCapacity.unlimited": "Used {{used}} / Unlimited", - "knowledgeBase.personalCapacity.available": "Available", - "knowledgeBase.personalCapacity.total": "Total capacity", - "knowledgeBase.personalCapacity.unlimitedValue": "Unlimited", - "knowledgeBase.personalCapacity.loadFailed": "Personal knowledge base capacity is temporarily unavailable", - "knowledgeBase.capacity.title": "Storage capacity", - "knowledgeBase.capacity.available": "Available", - "knowledgeBase.capacity.total": "Total capacity", - "knowledgeBase.capacity.unlimited": "Unlimited", "knowledgeBase.button.create": "Create", "knowledgeBase.button.sync": "Sync", "knowledgeBase.button.syncDataMate": "Sync DataMate Knowledge Bases", @@ -897,11 +867,6 @@ "knowledgeBase.search.placeholder": "Search knowledge base name", "knowledgeBase.filter.source.placeholder": "Filter by source", "knowledgeBase.filter.model.placeholder": "Filter by model", - "knowledgeBase.filter.title": "Filter knowledge bases", - "knowledgeBase.filter.button": "Filter", - "knowledgeBase.card.create": "Create knowledge base", - "knowledgeBase.card.createDescription": "Upload documents to build a knowledge asset", - "knowledgeBase.tag.updatedAt": "Updated {{date}}", "knowledgeBase.filter.clear": "Clear filters", "knowledgeBase.source.nexent": "{productName}", "knowledgeBase.source.datamate": "DataMate", @@ -998,7 +963,6 @@ "document.hint.uploadToCreate": "Please select files to upload to complete knowledge base creation", "document.hint.noDocuments": "No documents in this knowledge base, please upload documents", "document.table.header.name": "Document Name", - "document.table.header.tags": "Tags", "document.table.header.status": "Status", "document.table.header.size": "Size", "document.table.header.date": "Upload Date", @@ -4005,6 +3969,8 @@ "aidpKnowledge.kbUnavailable": "Unavailable", "aidpKnowledge.createKb": "Create Knowledge Base", "aidpKnowledge.refresh": "Refresh", + "aidpKnowledge.searchPlaceholder": "Search knowledge bases", + "aidpKnowledge.searchEmpty": "No matching knowledge bases", "aidpKnowledge.listEmpty": "No knowledge bases found", "aidpKnowledge.tagDocs": "docs: {{count}}", "aidpKnowledge.tagChunks": "chunks: {{count}}", @@ -4045,15 +4011,6 @@ "aidpKnowledge.docStatusFailed": "Failed", "aidpKnowledge.docSize": "Size", "aidpKnowledge.docCreatedAt": "Created At", - "aidpKnowledge.docActions": "Actions", - "aidpKnowledge.download": "Download", - "aidpKnowledge.delete": "Delete", - "aidpKnowledge.downloadSuccess": "File download started", - "aidpKnowledge.downloadFailed": "Failed to download file", - "aidpKnowledge.confirmDeleteDocTitle": "Delete document", - "aidpKnowledge.confirmDeleteDocContent": "Are you sure you want to delete this document? This cannot be undone.", - "aidpKnowledge.deleteDocSuccess": "Document deleted successfully", - "aidpKnowledge.deleteDocFailed": "Failed to delete document", "aidpKnowledge.noDocuments": "No documents yet", "aidpKnowledge.loadingDocs": "Loading documents...", "aidpKnowledge.uploadSuccess": "{{count}} document(s) uploaded successfully", diff --git a/frontend/public/locales/zh/common.json b/frontend/public/locales/zh/common.json index fe9c0613f..3f7d9f932 100644 --- a/frontend/public/locales/zh/common.json +++ b/frontend/public/locales/zh/common.json @@ -814,37 +814,7 @@ "knowledgeBase.error.syncFailed": "同步 DataMate 知识库失败", "knowledgeBase.message.testingConnection": "正在测试连接...", "knowledgeBase.message.testingSync": "正在同步知识库...", - "knowledgeBase.page.title": "知识库", - "knowledgeBase.page.description": "创建、管理并维护团队的知识资产", - "knowledgeBase.page.all": "全部知识库", - "knowledgeBase.page.count": "共 {{count}} 个", - "knowledgeBase.page.back": "返回知识库", - "knowledgeBase.create.subtitle": "配置基本信息,后续可随时调整", - "knowledgeBase.create.field.name": "知识库名称", - "knowledgeBase.create.field.description": "描述", - "knowledgeBase.create.optional": "(选填)", - "knowledgeBase.create.namePlaceholder": "例如:产品知识中心", - "knowledgeBase.create.descriptionPlaceholder": "简要说明这个知识库包含什么内容", - "knowledgeBase.create.advancedSettings": "高级设置", - "knowledgeBase.create.submit": "创建并进入", - "knowledgeBase.create.field.embeddingModel": "向量模型", - "knowledgeBase.create.field.groups": "所属用户组", - "knowledgeBase.create.field.permission": "组内权限", - "knowledgeBase.create.field.preserve": "文档副本", - "knowledgeBase.create.field.quota": "存储配额", - "knowledgeBase.create.uploadTitle": "上传文档,构建知识库", "knowledgeBase.list.title": "知识库列表", - "knowledgeBase.personalCapacity.title": "个人知识库容量", - "knowledgeBase.personalCapacity.withQuota": "已用 {{used}} / 总容量 {{quota}}", - "knowledgeBase.personalCapacity.unlimited": "已用 {{used}} / 无限制", - "knowledgeBase.personalCapacity.available": "可用容量", - "knowledgeBase.personalCapacity.total": "总容量", - "knowledgeBase.personalCapacity.unlimitedValue": "无限制", - "knowledgeBase.personalCapacity.loadFailed": "个人知识库容量暂时无法加载", - "knowledgeBase.capacity.title": "存储容量", - "knowledgeBase.capacity.available": "可用容量", - "knowledgeBase.capacity.total": "总容量", - "knowledgeBase.capacity.unlimited": "无限制", "knowledgeBase.button.create": "创建", "knowledgeBase.button.sync": "同步", "knowledgeBase.button.syncDataMate": "同步DataMate知识库", @@ -867,11 +837,6 @@ "knowledgeBase.search.placeholder": "搜索知识库名称", "knowledgeBase.filter.source.placeholder": "筛选来源", "knowledgeBase.filter.model.placeholder": "筛选模型", - "knowledgeBase.filter.title": "筛选知识库", - "knowledgeBase.filter.button": "筛选", - "knowledgeBase.card.create": "新建知识库", - "knowledgeBase.card.createDescription": "上传文档,构建专属知识资产", - "knowledgeBase.tag.updatedAt": "更新于{{date}}", "knowledgeBase.source.nexent": "{productName}", "knowledgeBase.source.datamate": "DataMate", "knowledgeBase.source.dify": "Dify", @@ -966,7 +931,6 @@ "document.hint.uploadToCreate": "请选择文件上传以完成知识库创建", "document.hint.noDocuments": "该知识库中暂无文档,请上传文档", "document.table.header.name": "文档名称", - "document.table.header.tags": "标签", "document.table.header.status": "状态", "document.table.header.size": "大小", "document.table.header.date": "上传日期", @@ -3720,6 +3684,8 @@ "aidpKnowledge.kbUnavailable": "不可用", "aidpKnowledge.createKb": "创建知识库", "aidpKnowledge.refresh": "刷新", + "aidpKnowledge.searchPlaceholder": "搜索知识库名称", + "aidpKnowledge.searchEmpty": "没有匹配的知识库", "aidpKnowledge.listEmpty": "暂无知识库", "aidpKnowledge.tagDocs": "文档: {{count}}", "aidpKnowledge.tagChunks": "分块: {{count}}", @@ -3760,15 +3726,6 @@ "aidpKnowledge.docStatusFailed": "失败", "aidpKnowledge.docSize": "大小", "aidpKnowledge.docCreatedAt": "创建时间", - "aidpKnowledge.docActions": "操作", - "aidpKnowledge.download": "下载", - "aidpKnowledge.delete": "删除", - "aidpKnowledge.downloadSuccess": "文件下载已开始", - "aidpKnowledge.downloadFailed": "文件下载失败", - "aidpKnowledge.confirmDeleteDocTitle": "删除文档", - "aidpKnowledge.confirmDeleteDocContent": "确定删除该文档吗?此操作无法撤销。", - "aidpKnowledge.deleteDocSuccess": "文档删除成功", - "aidpKnowledge.deleteDocFailed": "文档删除失败", "aidpKnowledge.noDocuments": "暂无文档", "aidpKnowledge.loadingDocs": "正在加载文档...", "aidpKnowledge.uploadSuccess": "成功上传 {{count}} 个文档", diff --git a/frontend/services/api.ts b/frontend/services/api.ts index b1c133628..a6fa45e0c 100644 --- a/frontend/services/api.ts +++ b/frontend/services/api.ts @@ -377,10 +377,6 @@ export const API_ENDPOINTS = { kbDetail: (id: string) => `${API_BASE_URL}/aidp-mgmt/knowledge-bases/${id}`, kbDocuments: (id: string) => `${API_BASE_URL}/aidp-mgmt/knowledge-bases/${id}/documents`, - removeKbDocuments: (id: string) => - `${API_BASE_URL}/aidp-mgmt/knowledge-bases/${id}/documents/remove`, - downloadKbDocument: (id: string) => - `${API_BASE_URL}/aidp-mgmt/knowledge-bases/${id}/documents/download`, models: `${API_BASE_URL}/aidp-mgmt/models`, /** PATCH endpoint for the per-KB in-group permission. */ kbPermission: (id: string) => diff --git a/sdk/nexent/scheduler/core.py b/sdk/nexent/scheduler/core.py index f7377a906..228730bf2 100644 --- a/sdk/nexent/scheduler/core.py +++ b/sdk/nexent/scheduler/core.py @@ -109,6 +109,7 @@ def __init__( self._loop_task: asyncio.Task[None] | None = None self._stop_event = asyncio.Event() self._running: set[asyncio.Task[None]] = set() + self._waiting: set[Hashable] = set() self._recovery_pending = True @property @@ -119,6 +120,25 @@ def is_running(self) -> bool: def active_count(self) -> int: return len(self._running) + @property + def waiting_count(self) -> int: + """Jobs parked on external (e.g. human) input; they still hold leases.""" + return len(self._waiting) + + def mark_waiting(self, job_id: Hashable, *, waiting: bool) -> None: + """Flag a claimed job as parked on external input (human decisions). + + A waiting job keeps its executor task and lease alive but no longer + consumes a max_concurrency slot, so hour-long human waits cannot + starve machine execution. The waiting set is bounded only by job + expiration, not by max_concurrency. Must be called on the scheduler + loop thread; worker threads should relay via loop.call_soon_threadsafe. + """ + if waiting: + self._waiting.add(job_id) + else: + self._waiting.discard(job_id) + async def start(self) -> None: if self.is_running: return @@ -153,7 +173,9 @@ async def _run_loop(self) -> None: if self._recovery_pending: await self.store.recover() self._recovery_pending = False - capacity = max(0, self.config.max_concurrency - len(self._running)) + # Waiting jobs still hold executor tasks, but their parked + # human-input waits must not consume execution concurrency. + capacity = max(0, self.config.max_concurrency - (len(self._running) - len(self._waiting))) if capacity: claimed = await self.store.claim_due( self.owner_id, @@ -203,6 +225,7 @@ async def _run_claimed(self, job: ClaimedJob[JobPayload]) -> None: except Exception: logger.exception("Scheduled job failed: job_id=%s", job.job_id) finally: + self._waiting.discard(job.job_id) renewal.cancel() await asyncio.gather(renewal, return_exceptions=True) try: diff --git a/test/backend/services/test_application_execute_attempt.py b/test/backend/services/test_application_execute_attempt.py index cff6064a2..ed4f462b8 100644 --- a/test/backend/services/test_application_execute_attempt.py +++ b/test/backend/services/test_application_execute_attempt.py @@ -1,7 +1,7 @@ """Unit test for application.py execute_attempt async consumer loop. -Covers the _flush_if_due batch/threshold flush (peek-then-take, begin_emit / -end_emit lifecycle) and the final flush in the ``finally`` block. No real +Covers the _flush_if_due batch/threshold flush (peek-then-drain_and_emit) and +the final flush in the ``finally`` block. No real database or agent — every collaborator is mocked. This lets us exercise the patch lines added in the HITL reorder PR without depending on the full Postgres-backed test suite. @@ -43,7 +43,7 @@ class _Port(RuntimeInteractionPort): def __init__(self, *a, **kw): self._chunk_buffer = [] self._chunk_buffer_lock = threading.Lock() - self._emit_in_flight = threading.Event() + self._emit_lock = threading.Lock() self.service = MagicMock() self.run_id = "run-1" self.tenant_id = "tenant-1" @@ -169,17 +169,12 @@ async def prepare_mock(**kwargs): return fake_info, None with _patched_application(make_port, fake_stream, prepare_mock) as application: - try: - await application.execute_attempt(*_execute_attempt_args()) - except StopAsyncIteration: - # The async consumer loop hit the StopAsyncIteration re-raised - # from the finally block — that's fine, we still ran through - # the entire consumer loop including _flush_if_due + final flush. - pass - except Exception as exc: - pytest.fail(f"execute_attempt raised unexpected {type(exc).__name__}: {exc!r}") + # A normally exhausted stream must NOT leak StopAsyncIteration — it + # used to escape the finally block and fail the claiming scheduler job. + await application.execute_attempt(*_execute_attempt_args()) port = port_ref["p"] + assert calls == [("finish", "failed")], calls total_persisted = sum(len(batch) for batch in port._emits) # Timing-sensitive — the sleep in fake_stream may cause the "final" chunk # to land in either the timed _flush_if_due path or the final-flush path. @@ -188,9 +183,9 @@ async def prepare_mock(**kwargs): f"Expected >=3 chunks persisted, got {total_persisted} batches={port._emits}" ) - # begin_emit must have been paired with end_emit — otherwise we would - # have seen _emit_in_flight still set after the loop. - assert not port._emit_in_flight.is_set(), "begin_emit without end_emit leaked" + # Every drain must have released the emit lock — a leak would deadlock + # the worker's flush before its next HITL transaction. + assert not port._emit_lock.locked(), "drain_and_emit leaked the emit lock" # --- exception / finally path coverage --------------------------------------- @@ -264,7 +259,7 @@ def add_chunk(chunk): # The two buffered chunks were flushed by the finally block, before finish. assert port._emits == [["chunk-0", "chunk-1"]], port._emits assert refs["finish_calls"] == ["failed"] - assert not port._emit_in_flight.is_set(), "begin_emit without end_emit leaked" + assert not port._emit_lock.locked(), "drain_and_emit leaked the emit lock" async def test_cancelled_error_cancels_scope_and_reraises_without_finish(): diff --git a/test/backend/services/test_runtime_port_chunk_buffer.py b/test/backend/services/test_runtime_port_chunk_buffer.py index e7eb40f08..168485a3e 100644 --- a/test/backend/services/test_runtime_port_chunk_buffer.py +++ b/test/backend/services/test_runtime_port_chunk_buffer.py @@ -9,11 +9,11 @@ import threading import time -import types from contextlib import contextmanager from unittest.mock import MagicMock, patch import pytest +from nexent.core.human_interaction.contracts import AttemptSuspended def _make_port(): @@ -199,7 +199,6 @@ def test_dispatch_calls_flush_before_transaction(): # Interaction path needs payload with questions for the early-return # "already answered" branches — let's use no interaction so we fall # through to the "not suspended → emit STARTED" case. - from nexent.core.human_interaction.contracts import AttemptSuspended # We expect either normal return or AttemptSuspended; both are fine as # long as flush_chunks_until_idle is called before any DB work. @@ -291,17 +290,17 @@ def fake_repo_transaction(*args, **kwargs): # --- in-flight emit and failure recovery -------------------------------------- def test_flush_until_idle_treats_in_flight_emit_as_busy(): - """An empty buffer is NOT idle while an async emit is handing chunks to - the DB thread — flush keeps polling until the emit completes. + """An empty buffer is NOT idle while an async drain holds the emit lock — + flush keeps polling until the drain completes. """ port = _make_port() port.emit_chunks = MagicMock() - port._emit_in_flight.set() + port._emit_lock.acquire() result: dict = {} def clear_soon(): time.sleep(0.06) - port._emit_in_flight.clear() + port._emit_lock.release() threading.Thread(target=clear_soon, daemon=True).start() @@ -318,7 +317,7 @@ def run_flush(): # Without the in-flight guard the flush would settle at ~10ms on the # empty buffer; observing >= 50ms proves it waited for the emit. assert result["elapsed"] >= 0.05, result - assert not port._emit_in_flight.is_set() + assert not port._emit_lock.locked() def test_flush_until_idle_restores_chunks_when_emit_chunks_raises(): @@ -334,7 +333,116 @@ def test_flush_until_idle_restores_chunks_when_emit_chunks_raises(): port.flush_chunks_until_idle(max_wait_ms=200, settle_ms=10) assert port.peek_chunks() == 2 - assert not port._emit_in_flight.is_set(), "begin_emit without end_emit leaked" + assert not port._emit_lock.locked(), "emit lock leaked on failure" + + +# --- SSE ordering race regressions -------------------------------------------- + +def test_flush_waits_for_in_flight_drain_past_deadline(): + """Regression: the hard deadline must NOT fire while an async drain holds + the emit lock, even past max_wait_ms. + + On a loaded server the async drain's run_blocking(emit_chunks) can exceed + 500ms; the old implementation broke out of the flush on the deadline and + wrote the HITL row, so the drained chunks landed after it in DB seq order + (the "form before model output" disorder). The flush must instead wait + for the lock, drain what arrived, and only then return. + """ + port = _make_port() + port._emit_lock.acquire() + result: dict = {} + + def drain_slowly(): + # Simulate an in-flight async drain that finishes after the deadline. + time.sleep(0.6) + port.add_chunk("late-chunk") + port._emit_lock.release() + + threading.Thread(target=drain_slowly, daemon=True).start() + + def run_flush(): + started = time.monotonic() + port.flush_chunks_until_idle(max_wait_ms=100, settle_ms=10) + result["elapsed"] = time.monotonic() - started + + worker = threading.Thread(target=run_flush, daemon=True) + worker.start() + worker.join(timeout=5) + + assert not worker.is_alive(), "flush blocked past the join timeout" + # The flush must have waited for the in-flight drain (0.6s), well past + # its own 100ms deadline, so the late chunks are persisted BEFORE the + # caller proceeds to write its HITL transaction. + assert result["elapsed"] >= 0.6, result + assert port.emit_chunks.call_args.args[0] == ["late-chunk"] + assert not port._emit_lock.locked() + + +def test_concurrent_drain_and_flush_preserve_chunk_order(): + """Regression: take_chunks + emit_chunks must be one atomic critical section. + + With the old design (take outside the lock) a worker flush and the async + drain could drain concurrently and the DB seq assignment followed lock + acquisition order instead of production order. With the shared critical + section, a chunk added earlier can never be persisted later. + """ + + def run_round(): + port = _make_port() + emitted: list[str] = [] + emit_lock = threading.Lock() + + def record(chunks, *, _lock=emit_lock, _out=emitted): + with _lock: + _out.extend(chunks) + + port.emit_chunks = MagicMock(side_effect=record) + total = 100 + + def drainer(stop: threading.Event, *, _port=port): + while not stop.is_set() or _port.peek_chunks(): + _port.drain_and_emit() + + stop = threading.Event() + threads = [threading.Thread(target=drainer, args=(stop,), daemon=True) for _ in range(2)] + for t in threads: + t.start() + for i in range(total): + port.add_chunk(f"c{i}") + stop.set() + for t in threads: + t.join(timeout=5) + assert not t.is_alive(), "drainer thread did not finish" + + assert emitted == [f"c{i}" for i in range(total)], ( + f"chunk order corrupted: {emitted[:10]}..." + ) + assert port.peek_chunks() == 0 + assert not port._emit_lock.locked() + + for _round in range(10): + run_round() + + +def test_drain_and_emit_empty_buffer_is_noop(): + """drain_and_emit on an empty buffer must not call emit_chunks.""" + port = _make_port() + port.drain_and_emit() + port.emit_chunks.assert_not_called() + assert not port._emit_lock.locked() + + +def test_drain_and_emit_restores_chunks_on_failure(): + """A failed DB emit inside drain_and_emit puts the chunks back.""" + port = _make_port() + port.emit_chunks = MagicMock(side_effect=RuntimeError("db down")) + port.add_chunk("x") + + with pytest.raises(RuntimeError, match="db down"): + port.drain_and_emit() + + assert port.peek_chunks() == 1 + assert not port._emit_lock.locked() def test_finish_flush_failure_still_writes_terminal_status(): diff --git a/test/ext_components/aidp/mock_servers/aidp_mgmt_mock_server.py b/test/ext_components/aidp/mock_servers/aidp_mgmt_mock_server.py index 676ae6e7a..ae5c86474 100644 --- a/test/ext_components/aidp/mock_servers/aidp_mgmt_mock_server.py +++ b/test/ext_components/aidp/mock_servers/aidp_mgmt_mock_server.py @@ -9,8 +9,6 @@ - DELETE /KnowledgeBase/Tenants/{tenant}/KnowledgeBases/{id} (delete) - POST /KnowledgeBase/Tenants/{tenant}/KnowledgeBases/{id}/KnowledgeFiles/Upload (upload docs) - GET /KnowledgeBase/Tenants/{tenant}/KnowledgeBases/{id}/KnowledgeFiles (list docs) - - POST /KnowledgeBase/Tenants/{tenant}/KnowledgeBases/{id}/KnowledgeFiles/Remove (remove docs) - - POST /KnowledgeBase/Tenants/{tenant}/KnowledgeBases/{id}/KnowledgeFiles/Download (download doc) - GET /KnowledgeBase/Tenants/{tenant}/KnowledgeBases/{id}/Channels (ingestion channels) - POST /KnowledgeBase/Tenants/{tenant}/KnowledgeBases/{id}/KnowledgeFiles/History (all-status file history) - POST /KnowledgeBase/Tenants/{tenant}/Retrieval/FusionSearch (search - preserved from reference) @@ -36,18 +34,16 @@ import argparse import json import logging -import mimetypes import os import time import uuid from pathlib import Path from typing import Any, Dict, List, Literal, Optional -from urllib.parse import quote from fastapi import FastAPI, File, Header, HTTPException, Query, UploadFile from fastapi.middleware.cors import CORSMiddleware from pydantic import BaseModel, Field -from starlette.responses import JSONResponse, StreamingResponse +from starlette.responses import JSONResponse logger = logging.getLogger("aidp_mgmt_mock") logging.basicConfig( @@ -151,16 +147,14 @@ def _seed_initial_data() -> None: # Seed some documents for the FAQ KB so list_docs is non-empty by default. _DOCUMENTS_BY_KB["aidp-kb-faq"] = [ { - "file_uuid": "00000000-0000-4000-8000-000000000001", - "file_ino_no": 1001, + "file_ino_no": "file-faq-001", "file_name": "常见问题汇总.txt", "file_size": 2048, "file_type": "txt", "create_time": 1718000400, }, { - "file_uuid": "00000000-0000-4000-8000-000000000002", - "file_ino_no": 1002, + "file_ino_no": "file-faq-002", "file_name": "troubleshooting.md", "file_size": 4096, "file_type": "md", @@ -227,55 +221,7 @@ def _load_state() -> None: ) -def _ensure_document_uuids() -> None: - """Backfill stable UUIDs for state created before UUID support existed.""" - for kds_id, documents in _DOCUMENTS_BY_KB.items(): - if not isinstance(documents, list): - continue - for document in documents: - if not isinstance(document, dict) or document.get("file_uuid"): - continue - file_ino_no = str(document.get("file_ino_no") or uuid.uuid4()) - document["file_uuid"] = str( - uuid.uuid5(uuid.NAMESPACE_URL, f"mock-aidp:{kds_id}:{file_ino_no}") - ) - - _load_state() -_ensure_document_uuids() -_save_state() - - -def _public_document(document: Dict[str, Any]) -> Dict[str, Any]: - """Return document metadata without any mock-only private fields.""" - return {key: value for key, value in document.items() if not key.startswith("_")} - - -def _find_document(kds_id: str, file_uuid: str) -> Optional[Dict[str, Any]]: - return next( - ( - document - for document in _DOCUMENTS_BY_KB.get(kds_id, []) - if document.get("file_uuid") == file_uuid - ), - None, - ) - - -def _document_content(document: Dict[str, Any]) -> bytes: - """Build deterministic mock content for a document download.""" - return ( - f"Mock AIDP content for {document.get('file_name', 'download')}\n" - ).encode("utf-8") - - -def _content_disposition(filename: str) -> str: - """Build a standard ASCII fallback plus RFC 5987 UTF-8 filename header.""" - ascii_name = "".join( - char if 32 <= ord(char) < 127 and char not in {'"', "\\"} else "_" - for char in filename - ) - return f'attachment; filename="{ascii_name}"; filename*=UTF-8\'\'{quote(filename)}' # ============================================================================= @@ -374,14 +320,6 @@ class FusionSearchRequest(BaseModel): metadata_condition: Optional[MetadataCondition] = None -class RemoveFilesBody(BaseModel): - file_uuids: List[uuid.UUID] = Field(..., min_length=1) - - -class DownloadFileBody(BaseModel): - file_uuid: uuid.UUID = Field(...) - - # ============================================================================= # Auth helper # ============================================================================= @@ -711,16 +649,8 @@ async def upload_documents( for f in files: try: content = await f.read() - file_ino_no = max( - ( - document["file_ino_no"] - for document in _DOCUMENTS_BY_KB.get(kds_id, []) - if isinstance(document.get("file_ino_no"), int) - ), - default=0, - ) + 1 + file_ino_no = f"file-{uuid.uuid4().hex[:12]}" doc = { - "file_uuid": str(uuid.uuid4()), "file_ino_no": file_ino_no, "file_name": f.filename or "unknown", "file_size": len(content), @@ -756,7 +686,6 @@ async def upload_documents( "file_type": doc["file_type"], "file_size": doc["file_size"], "file_ino_no": doc["file_ino_no"], - "file_uuid": doc["file_uuid"], "first_upload_time": doc["create_time"], } for doc in success_docs @@ -789,7 +718,7 @@ def list_documents( ] start = (page - 1) * page_size end = start + page_size - items = [_public_document(doc) for doc in all_docs[start:end]] + items = all_docs[start:end] # Real AIDP returns `next_link` as the authoritative "more pages exist" # signal. When there are no more docs, next_link is simply absent. @@ -806,89 +735,6 @@ def list_documents( }) -@app.post(f"{_KB_PREFIX}/{{kds_id}}/KnowledgeFiles/Remove") -def remove_documents( - kds_id: str, - body: RemoveFilesBody, - authorization: Optional[str] = Header(default=None), -) -> JSONResponse: - """Remove documents by file UUID and return per-file success/failure lists.""" - _check_auth(authorization) - - if kds_id not in _KNOWLEDGE_BASES: - raise HTTPException(status_code=404, detail=f"Knowledge base {kds_id} not found") - - documents = _DOCUMENTS_BY_KB.setdefault(kds_id, []) - remaining = list(documents) - success_list: List[Dict[str, str]] = [] - failed_list: List[Dict[str, str]] = [] - for raw_file_uuid in body.file_uuids: - file_uuid = str(raw_file_uuid) - matched = next( - (document for document in remaining if document.get("file_uuid") == file_uuid), - None, - ) - if matched is None: - failed_list.append({"file_uuid": file_uuid}) - continue - remaining.remove(matched) - success_list.append({"file_uuid": file_uuid}) - - _DOCUMENTS_BY_KB[kds_id] = remaining - _save_state() - logger.info( - "REMOVE DOCS kds_id=%s total=%d success=%d failed=%d", - kds_id, - len(body.file_uuids), - len(success_list), - len(failed_list), - ) - return JSONResponse(content={ - "summary": { - "total": len(body.file_uuids), - "success": len(success_list), - "failed": len(failed_list), - }, - "success_list": success_list, - "failed_list": failed_list, - }) - - -@app.post(f"{_KB_PREFIX}/{{kds_id}}/KnowledgeFiles/Download") -def download_document( - kds_id: str, - body: DownloadFileBody, - authorization: Optional[str] = Header(default=None), -) -> StreamingResponse: - """Return deterministic binary content for a document download.""" - _check_auth(authorization) - - if kds_id not in _KNOWLEDGE_BASES: - raise HTTPException(status_code=404, detail=f"Knowledge base {kds_id} not found") - - file_uuid = str(body.file_uuid) - document = _find_document(kds_id, file_uuid) - if document is None: - raise HTTPException(status_code=404, detail=f"File {file_uuid} not found") - - filename = str(document.get("file_name") or "download") - content = _document_content(document) - content_type = mimetypes.guess_type(filename)[0] or "application/octet-stream" - response_headers = { - "Content-Disposition": _content_disposition(filename), - "X-File-Size": str(len(content)), - } - async def content_stream(): - for offset in range(0, len(content), 8 * 1024): - yield content[offset : offset + 8 * 1024] - - return StreamingResponse( - content_stream(), - media_type=content_type, - headers=response_headers, - ) - - # ============================================================================= # Ingestion channels + knowledge-file history # ============================================================================= diff --git a/test/ext_components/aidp/test_aidp_mgmt_app.py b/test/ext_components/aidp/test_aidp_mgmt_app.py index bb845e2cb..9d5c6305d 100644 --- a/test/ext_components/aidp/test_aidp_mgmt_app.py +++ b/test/ext_components/aidp/test_aidp_mgmt_app.py @@ -21,7 +21,6 @@ from typing import Any from unittest.mock import MagicMock, patch -import httpx import pytest from fastapi import FastAPI from fastapi.testclient import TestClient @@ -47,11 +46,6 @@ def _mod(name): nexent_storage_factory = _mod("nexent.storage.storage_client_factory") nexent_storage_factory.create_storage_client_from_config = MagicMock() -services_pkg = _mod("services") -services_pkg.__path__ = [os.path.join(BACKEND_DIR, "services")] -tag_management_service = _mod("services.tag_management_service") -tag_management_service.TagManagementService = MagicMock() - class _MinIOStorageConfig: def __init__(self, **kwargs): @@ -61,7 +55,7 @@ def __init__(self, **kwargs): nexent_storage_factory.MinIOStorageConfig = _MinIOStorageConfig for mod in (nexent_pkg, nexent_utils, nexent_http_mgr, nexent_storage, - nexent_storage_factory, services_pkg, tag_management_service): + nexent_storage_factory): sys.modules.setdefault(mod.__name__, mod) # Register non-prefixed ``database`` / ``database.client`` stubs so that @@ -854,153 +848,6 @@ def count_docs(*_args, **_kwargs): assert response.json()["total_count"] == 1 -# --- Remove/download documents ------------------------------------------- - - -class TestAidpDocumentFileOperations: - def test_remove_forwards_only_uuids(self): - client = _client() - from ext_components.aidp.apps import aidp_mgmt_app - from ext_components.aidp.services import aidp_permission_service - - aidp_result = { - "summary": {"total": 2, "success": 1, "failed": 1}, - "success_list": [{"file_uuid": "00000000-0000-4000-8000-000000000001"}], - "failed_list": [{"file_uuid": "00000000-0000-4000-8000-000000000002"}], - } - with patch.object( - aidp_permission_service, - "require_permission", - return_value=MagicMock(permission="EDIT"), - ), patch.object( - aidp_mgmt_app, - "remove_aidp_docs_impl", - return_value=aidp_result, - ) as mock_remove: - response = client.post( - "/aidp-mgmt/knowledge-bases/kb-1/documents/remove", - headers=_bearer(), - json={ - "file_uuids": [ - "00000000-0000-4000-8000-000000000001", - "00000000-0000-4000-8000-000000000002", - ] - }, - ) - - assert response.status_code == HTTPStatus.OK - assert response.json() == aidp_result - assert mock_remove.call_args.args[3] == [ - "00000000-0000-4000-8000-000000000001", - "00000000-0000-4000-8000-000000000002", - ] - - def test_remove_requires_file_uuids(self): - client = _client() - from ext_components.aidp.apps import aidp_mgmt_app - - with patch.object(aidp_mgmt_app, "remove_aidp_docs_impl") as mock_remove: - response = client.post( - "/aidp-mgmt/knowledge-bases/kb-1/documents/remove", - headers=_bearer(), - json={}, - ) - - assert response.status_code == HTTPStatus.UNPROCESSABLE_ENTITY - mock_remove.assert_not_called() - - def test_remove_requires_standard_file_uuid(self): - client = _client() - from ext_components.aidp.apps import aidp_mgmt_app - - with patch.object(aidp_mgmt_app, "remove_aidp_docs_impl") as mock_remove: - response = client.post( - "/aidp-mgmt/knowledge-bases/kb-1/documents/remove", - headers=_bearer(), - json={"file_uuids": ["uuid-1"]}, - ) - - assert response.status_code == HTTPStatus.UNPROCESSABLE_ENTITY - mock_remove.assert_not_called() - - def test_download_returns_binary_response(self): - client = _client() - from ext_components.aidp.apps import aidp_mgmt_app - from ext_components.aidp.services import aidp_permission_service - - mock_response = httpx.Response( - 200, - headers={ - "Content-Type": "text/plain", - "Content-Disposition": 'attachment; filename="a.txt"', - "X-File-Size": "5", - }, - content=b"hello", - request=httpx.Request("POST", SERVER_URL), - ) - - async def stream_document(*args): - return mock_response - - with patch.object( - aidp_permission_service, - "require_permission", - return_value=MagicMock(permission="READ_ONLY"), - ), patch.object( - aidp_mgmt_app, - "stream_aidp_doc_impl", - side_effect=stream_document, - ) as mock_download: - response = client.post( - "/aidp-mgmt/knowledge-bases/kb-1/documents/download", - headers=_bearer(), - json={"file_uuid": "00000000-0000-4000-8000-000000000001"}, - ) - - assert response.status_code == HTTPStatus.OK - assert response.content == b"hello" - assert response.headers["content-type"].startswith("text/plain") - assert response.headers["content-disposition"] == 'attachment; filename="a.txt"' - assert response.headers["x-file-size"] == "5" - assert mock_download.call_args.args == ( - SERVER_URL, - API_KEY, - "kb-1", - "00000000-0000-4000-8000-000000000001", - ) - assert "x-file-name" not in response.headers - - def test_remove_does_not_invalidate_cache_when_no_file_succeeds(self): - client = _client() - from ext_components.aidp.apps import aidp_mgmt_app - from ext_components.aidp.services import aidp_permission_service - - aidp_result = { - "summary": {"total": 1, "success": 0, "failed": 1}, - "success_list": [], - "failed_list": [{"file_uuid": "00000000-0000-4000-8000-000000000001"}], - } - with patch.object( - aidp_permission_service, - "require_permission", - return_value=MagicMock(permission="EDIT"), - ), patch.object( - aidp_mgmt_app, - "remove_aidp_docs_impl", - return_value=aidp_result, - ), patch.object(aidp_mgmt_app, "invalidate_aidp_kb_detail_cache") as mock_kb_cache, patch.object( - aidp_mgmt_app, "invalidate_aidp_doc_count_cache" - ) as mock_count_cache: - response = client.post( - "/aidp-mgmt/knowledge-bases/kb-1/documents/remove", - headers=_bearer(), - json={"file_uuids": ["00000000-0000-4000-8000-000000000001"]}, - ) - - assert response.status_code == HTTPStatus.OK - mock_kb_cache.assert_not_called() - mock_count_cache.assert_not_called() - # --- Models list (auth only, no per-KB permission) ------------------------ diff --git a/test/ext_components/aidp/test_aidp_service.py b/test/ext_components/aidp/test_aidp_service.py index 46decb832..4fe5d35fc 100644 --- a/test/ext_components/aidp/test_aidp_service.py +++ b/test/ext_components/aidp/test_aidp_service.py @@ -3,7 +3,7 @@ import os import sys from types import ModuleType -from unittest.mock import AsyncMock, MagicMock +from unittest.mock import MagicMock import httpx import pytest @@ -1201,16 +1201,6 @@ def _setup_mock_client(aidp_service_module, method="get", response=None, side_ef return mock_client -def _setup_mock_async_client(aidp_service_module, response=None, side_effect=None): - """Create and wire an async mock client into the service module manager.""" - mock_client = MagicMock() - mock_client.send = AsyncMock(side_effect=side_effect, return_value=response) - mock_manager = MagicMock() - mock_manager.get_async_client.return_value = mock_client - aidp_service_module.http_client_manager = mock_manager - return mock_client - - def _make_http_error(status_code, method="GET"): """Create an httpx.HTTPStatusError with given status code.""" request = httpx.Request(method, "http://127.0.0.1:30081") @@ -2115,7 +2105,7 @@ def test_invalid_config(self, aidp_service_module, server_url, api_key): def test_success_normalizes_docs(self, aidp_service_module): mock_resp = _make_success_response({ "value": [ - {"name": "doc1", "file_uuid": "uuid-1", "first_upload_time": 1700000000}, + {"name": "doc1", "first_upload_time": 1700000000}, {"name": "doc2", "create_time": 1700100000, "update_time": 1700200000}, ], "total_count": 2, @@ -2132,7 +2122,6 @@ def test_success_normalizes_docs(self, aidp_service_module): assert len(result["value"]) == 2 # Normalization adds created_at / updated_at assert result["value"][0]["created_at"] is not None - assert result["value"][0]["file_uuid"] == "uuid-1" assert result["value"][1]["updated_at"] is not None def test_success_non_list_value_not_normalized(self, aidp_service_module): @@ -2202,147 +2191,6 @@ def test_json_parse_value_error(self, aidp_service_module): assert exc_info.value.error_code == ErrorCode.AIDP_RESPONSE_ERROR -# --------------------------------------------------------------------------- -# remove_aidp_docs_impl / download_aidp_doc_impl tests -# --------------------------------------------------------------------------- -class TestAidpDocumentFileOperations: - def test_remove_sends_uuid_array_and_preserves_partial_result( - self, aidp_service_module - ): - expected = { - "summary": {"total": 2, "success": 1, "failed": 1}, - "success_list": [{"file_uuid": "uuid-1"}], - "failed_list": [{"file_uuid": "uuid-2"}], - } - mock_client = _setup_mock_client( - aidp_service_module, - method="post", - response=_make_success_response(expected), - ) - - result = aidp_service_module.remove_aidp_docs_impl( - "http://127.0.0.1:30081", - "jwt-token", - "kb-1", - ["uuid-1", "uuid-2"], - ) - - assert result == expected - call = mock_client.post.call_args - assert call.args[0].endswith( - "/KnowledgeBase/Tenants/aidp/KnowledgeBases/kb-1/KnowledgeFiles/Remove" - ) - assert call.kwargs["json"] == {"file_uuids": ["uuid-1", "uuid-2"]} - - def test_remove_maps_request_error(self, aidp_service_module): - request = httpx.Request("POST", "http://127.0.0.1:30081") - _setup_mock_client( - aidp_service_module, - method="post", - side_effect=httpx.RequestError("network down", request=request), - ) - - with pytest.raises(AppException) as exc_info: - aidp_service_module.remove_aidp_docs_impl( - "http://127.0.0.1:30081", "jwt-token", "kb-1", ["uuid-1"] - ) - assert exc_info.value.error_code == ErrorCode.AIDP_CONNECTION_ERROR - - def test_remove_maps_invalid_json(self, aidp_service_module): - mock_response = _make_success_response({}) - mock_response.json.side_effect = ValueError("bad json") - _setup_mock_client(aidp_service_module, method="post", response=mock_response) - - with pytest.raises(AppException) as exc_info: - aidp_service_module.remove_aidp_docs_impl( - "http://127.0.0.1:30081", "jwt-token", "kb-1", ["uuid-1"] - ) - assert exc_info.value.error_code == ErrorCode.AIDP_RESPONSE_ERROR - - @pytest.mark.parametrize("status_code", [401, 403, 500]) - def test_remove_maps_upstream_http_errors(self, aidp_service_module, status_code): - _setup_mock_client( - aidp_service_module, - method="post", - side_effect=_make_http_error(status_code, "POST"), - ) - - with pytest.raises(AppException) as exc_info: - aidp_service_module.remove_aidp_docs_impl( - "http://127.0.0.1:30081", "jwt-token", "kb-1", ["uuid-1"] - ) - expected_code = ( - ErrorCode.AIDP_AUTH_ERROR - if status_code in (401, 403) - else ErrorCode.AIDP_SERVICE_ERROR - ) - assert exc_info.value.error_code == expected_code - - @pytest.mark.asyncio - async def test_download_streams_binary_content_and_headers(self, aidp_service_module): - mock_response = httpx.Response( - 200, - headers={ - "Content-Type": "text/plain", - "Content-Disposition": 'attachment; filename="a.txt"', - "X-File-Size": "16", - }, - content=b"downloaded bytes", - request=httpx.Request("POST", "http://127.0.0.1:30081"), - ) - mock_client = _setup_mock_async_client(aidp_service_module, response=mock_response) - - response = await aidp_service_module.stream_aidp_doc_impl( - "http://127.0.0.1:30081", "jwt-token", "kb-1", "uuid-1" - ) - chunks = [chunk async for chunk in response.aiter_bytes()] - - assert b"".join(chunks) == b"downloaded bytes" - assert response.headers["Content-Type"] == "text/plain" - assert response.headers["Content-Disposition"] == 'attachment; filename="a.txt"' - assert response.headers["X-File-Size"] == "16" - request = mock_client.build_request.call_args - assert request.args[0] == "POST" - assert request.args[1].endswith( - "/KnowledgeBase/Tenants/aidp/KnowledgeBases/kb-1/KnowledgeFiles/Download" - ) - assert request.kwargs["json"] == {"file_uuid": "uuid-1"} - await response.aclose() - assert mock_response.is_closed - - @pytest.mark.asyncio - async def test_download_maps_request_error(self, aidp_service_module): - request = httpx.Request("POST", "http://127.0.0.1:30081") - _setup_mock_async_client( - aidp_service_module, - side_effect=httpx.RequestError("network down", request=request), - ) - - with pytest.raises(AppException) as exc_info: - await aidp_service_module.stream_aidp_doc_impl( - "http://127.0.0.1:30081", "jwt-token", "kb-1", "uuid-1" - ) - assert exc_info.value.error_code == ErrorCode.AIDP_CONNECTION_ERROR - - @pytest.mark.asyncio - async def test_download_maps_upstream_http_error(self, aidp_service_module): - response = httpx.Response( - 404, - json={"error": "file not found"}, - request=httpx.Request("POST", "http://127.0.0.1:30081"), - ) - _setup_mock_async_client( - aidp_service_module, - response=response, - ) - - with pytest.raises(AppException) as exc_info: - await aidp_service_module.stream_aidp_doc_impl( - "http://127.0.0.1:30081", "jwt-token", "kb-1", "uuid-1" - ) - assert exc_info.value.error_code == ErrorCode.AIDP_SERVICE_ERROR - - # --------------------------------------------------------------------------- # list_aidp_models_impl tests # --------------------------------------------------------------------------- diff --git a/test/sdk/scheduler/test_core.py b/test/sdk/scheduler/test_core.py index 4c51a65e6..e7994efe9 100644 --- a/test/sdk/scheduler/test_core.py +++ b/test/sdk/scheduler/test_core.py @@ -136,6 +136,45 @@ async def execute(job, lease): assert observed_lost == [True] +@pytest.mark.asyncio +async def test_waiting_jobs_do_not_consume_concurrency(): + store = MemoryLeaseStore([ClaimedJob(1, {"id": 1}), ClaimedJob(2, {"id": 2})]) + started = [] + first_running = asyncio.Event() + second_running = asyncio.Event() + gate = asyncio.Event() + + async def execute(job, lease): + started.append(job.job_id) + # Both executors must stay alive: an instantly-completing job 2 would + # race its done-callback against the active_count assertion below. + if job.job_id == 1: + first_running.set() + else: + second_running.set() + await gate.wait() + + scheduler = LeaseScheduler(store, execute, _config(max_concurrency=1), owner_id="scheduler-a") + await scheduler.start() + await _wait_until(first_running.is_set) + + # With job 1 executing, the single concurrency slot is fully consumed. + assert scheduler.active_count == 1 + assert 2 not in started + + # Job 1 parks on human input: its slot frees for job 2 without finishing. + scheduler.mark_waiting(1, waiting=True) + await _wait_until(second_running.is_set) + assert scheduler.active_count == 2 + assert scheduler.waiting_count == 1 + + gate.set() + await scheduler.stop() + + # Completion clears the waiting flag even without an explicit reset. + assert scheduler.waiting_count == 0 + + @pytest.mark.asyncio async def test_scheduler_recovers_after_transient_claim_failure(): store = MemoryLeaseStore([ClaimedJob(1, {"id": 1})])