diff --git a/crates/ov_cli/src/client.rs b/crates/ov_cli/src/client.rs index fe5a58d9a..26165efac 100644 --- a/crates/ov_cli/src/client.rs +++ b/crates/ov_cli/src/client.rs @@ -692,6 +692,7 @@ impl HttpClient { pub async fn add_resource( &self, path: &str, + add_type: Option, to: Option, parent: Option, parent_auto_create: Option, @@ -738,7 +739,9 @@ impl HttpClient { body }; - if path_obj.exists() { + // A declared Connector add_type sends the path verbatim as a remote + // source; never interpret it as a local file to upload. + if add_type.is_none() && path_obj.exists() { if path_obj.is_dir() { let source_name = path_obj .file_name() @@ -840,6 +843,7 @@ impl HttpClient { } else { let body = build_body(serde_json::json!({ "path": path, + "add_type": add_type, "to": to, "parent": effective_parent, "reason": reason, diff --git a/crates/ov_cli/src/commands/resources.rs b/crates/ov_cli/src/commands/resources.rs index 415eca5ff..0afdc2512 100644 --- a/crates/ov_cli/src/commands/resources.rs +++ b/crates/ov_cli/src/commands/resources.rs @@ -6,6 +6,7 @@ use serde_json::{Map, Value}; pub async fn add_resource( client: &HttpClient, path: &str, + add_type: Option, to: Option, parent: Option, parent_auto_create: Option, @@ -31,6 +32,7 @@ pub async fn add_resource( let result = client .add_resource( path, + add_type, to, parent, parent_auto_create, diff --git a/crates/ov_cli/src/handlers.rs b/crates/ov_cli/src/handlers.rs index 3d483e3d2..e5a5e4ff2 100644 --- a/crates/ov_cli/src/handlers.rs +++ b/crates/ov_cli/src/handlers.rs @@ -15,6 +15,7 @@ use serde_json::{Map, Value}; pub async fn handle_add_resource( mut path: String, + add_type: Option, to: Option, parent: Option, parent_auto_create: Option, @@ -37,7 +38,9 @@ pub async fn handle_add_resource( let is_url = path.starts_with("http://") || path.starts_with("https://") || path.starts_with("git@"); - if !is_url { + // A declared Connector add_type sends the path verbatim to the server; + // it is never a local file, so skip local-path existence validation. + if add_type.is_none() && !is_url { use std::path::Path; // Unescape path: replace backslash followed by space with just space @@ -98,6 +101,7 @@ pub async fn handle_add_resource( commands::resources::add_resource( &client, &path, + add_type, to, parent, parent_auto_create, diff --git a/crates/ov_cli/src/help_ui.rs b/crates/ov_cli/src/help_ui.rs index 182e68536..97d3a089b 100644 --- a/crates/ov_cli/src/help_ui.rs +++ b/crates/ov_cli/src/help_ui.rs @@ -142,6 +142,10 @@ const COMMAND_HELP_SPECS: &[CommandHelpSpec] = &[ label: "ov add-resource https://example.com/sitemap.xml --watch-interval 1440", description: "Import a whole site via sitemap/RSS and refresh it daily.", }, + HelpItem { + label: "ov add-resource tos://bucket/docs/ --add-type tos --to viking://resources/docs", + description: "Declare the Connector source type explicitly (Connector integration must be enabled).", + }, ], next_steps: &[ HelpItem { diff --git a/crates/ov_cli/src/main.rs b/crates/ov_cli/src/main.rs index ea8b7cf43..22017ef62 100644 --- a/crates/ov_cli/src/main.rs +++ b/crates/ov_cli/src/main.rs @@ -282,6 +282,18 @@ enum Commands { /// Local path or URL to import #[arg(value_name = "path-or-url")] path: String, + /// Explicit Connector source type (e.g. "tos", "git"). Routes the import + /// through the Connector integration (must be enabled server-side); the + /// path is sent verbatim and never treated as a local file. Requires --to + /// and cannot be combined with --parent or --parent-auto-create + #[arg( + long = "add-type", + value_name = "type", + requires = "to", + conflicts_with_all = ["parent", "parent_auto_create"], + help_heading = "Common options" + )] + add_type: Option, /// Exact target URI (must not exist yet) (cannot be used with --parent) #[arg(long, value_name = "uri", help_heading = "Common options")] to: Option, @@ -2798,6 +2810,7 @@ async fn main() { let result = match cli.command { Commands::AddResource { path, + add_type, to, parent, parent_auto_create, @@ -2821,6 +2834,7 @@ async fn main() { ctx.with_upload_options(upload_options.merged_with_legacy(legacy_upload_options)); handlers::handle_add_resource( path, + add_type, to, parent, parent_auto_create, @@ -4036,6 +4050,61 @@ mod tests { assert!(Cli::try_parse_from(["ov", "skills", "update", "--progress"]).is_err()); } + #[test] + fn cli_add_resource_add_type_requires_exact_to() { + assert!( + Cli::try_parse_from(["ov", "add-resource", "space:home", "--add-type", "feishu"]) + .is_err() + ); + assert!( + Cli::try_parse_from([ + "ov", + "add-resource", + "space:home", + "--add-type", + "feishu", + "--to", + "viking://resources/feishu", + "--parent", + "viking://resources/imports", + ]) + .is_err() + ); + assert!( + Cli::try_parse_from([ + "ov", + "add-resource", + "space:home", + "--add-type", + "feishu", + "--to", + "viking://resources/feishu", + "--parent-auto-create", + "viking://resources/imports", + ]) + .is_err() + ); + + let cli = Cli::try_parse_from([ + "ov", + "add-resource", + "space:home", + "--add-type", + "feishu", + "--to", + "viking://resources/feishu", + ]) + .expect("declared add type with an exact target should parse"); + + match cli.command { + Commands::AddResource { add_type, to, .. } => { + assert_eq!(add_type.as_deref(), Some("feishu")); + assert_eq!(to.as_deref(), Some("viking://resources/feishu")); + } + _ => panic!("expected add-resource command"), + } + } + #[test] fn cli_parses_add_resource_tags() { let cli = Cli::try_parse_from([ diff --git a/openviking/async_client.py b/openviking/async_client.py index f0c07d4d9..065f99093 100644 --- a/openviking/async_client.py +++ b/openviking/async_client.py @@ -330,6 +330,7 @@ class AsyncOpenViking: args: Optional[Dict[str, Any]] = None, telemetry: TelemetryRequest = False, processing_mode: str = "semantic_and_vectors", + add_type: Optional[str] = None, tags: Optional[List[str]] = None, tag_mode: str = "replace", **kwargs, @@ -341,6 +342,9 @@ class AsyncOpenViking: path: Local file path or URL. A sitemap / RSS / Atom URL ingests the whole site as one resource tree; pass ``args={"site": True}`` to force whole-site ingestion from a bare domain. + add_type: Explicit Connector source type. Requires an exact ``to`` + target and cannot be combined with ``parent``. The source + ``path`` is forwarded verbatim. reason: Context/reason for adding this resource. instruction: Specific instruction for processing. wait: If True, wait for processing to complete. @@ -357,11 +361,18 @@ class AsyncOpenViking: """ await self._ensure_initialized() + if add_type is not None: + add_type = add_type.strip() or None + if add_type and parent: + raise ValueError("'add_type' cannot be combined with 'parent'.") + if add_type and not to: + raise ValueError("'add_type' requires an exact 'to' target.") if to and parent: raise ValueError("Cannot specify both 'to' and 'parent' at the same time.") return await self._client.add_resource( path=path, + add_type=add_type, to=to, parent=parent, reason=reason, diff --git a/openviking/client/local.py b/openviking/client/local.py index abe51e488..40bdeb5db 100644 --- a/openviking/client/local.py +++ b/openviking/client/local.py @@ -139,11 +139,22 @@ class LocalClient(BaseClient): watch_interval: float = 0, args: Optional[Dict[str, Any]] = None, processing_mode: str = "semantic_and_vectors", + add_type: Optional[str] = None, tags: Optional[List[str]] = None, tag_mode: str = "replace", **kwargs, ) -> Dict[str, Any]: - """Add resource to OpenViking.""" + """Add resource to OpenViking. + + ``add_type`` declares a Connector source and requires an exact ``to`` + target; it cannot be combined with ``parent``. + """ + if add_type is not None: + add_type = add_type.strip() or None + if add_type and parent: + raise ValueError("'add_type' cannot be combined with 'parent'.") + if add_type and not to: + raise ValueError("'add_type' requires an exact 'to' target.") if to and parent: raise ValueError("Cannot specify both 'to' and 'parent' at the same time.") @@ -153,6 +164,7 @@ class LocalClient(BaseClient): fn=lambda: self._service.resources.add_resource( path=path, ctx=self._ctx, + add_type=add_type, to=to, parent=parent, reason=reason, diff --git a/openviking/connector/client.py b/openviking/connector/client.py index c0e8a9992..7da1a566b 100644 --- a/openviking/connector/client.py +++ b/openviking/connector/client.py @@ -70,8 +70,9 @@ class ConnectorClient: ``to`` is the exact OpenViking file or directory target. Source-specific settings stay inside ``param_config``; source credentials stay inside - ``auth_config``, the only body field redacted from request logs on - every hop, and must never be merged into ``param_config``. + ``auth_config``, which Connector and plugin request logs redact, and + must never be merged into ``param_config``. The incoming OpenViking + HTTP body may still be captured when unsafe body dumping is enabled. Returns the Connector response dict (contains task key / id on success). """ diff --git a/openviking/connector/delegate.py b/openviking/connector/delegate.py new file mode 100644 index 000000000..93eb56138 --- /dev/null +++ b/openviking/connector/delegate.py @@ -0,0 +1,612 @@ +# Copyright (c) 2026 Beijing Volcano Engine Technology Co., Ltd. +# SPDX-License-Identifier: AGPL-3.0 +"""Connector delegation for resource imports. + +ResourceService hands an add_resource request to :class:`ConnectorDelegate` +when the source belongs to the external Connector integration; the native +service keeps only that hand-off. This module owns the whole delegation +lifecycle: + +- the routing decision (:meth:`ConnectorDelegate.should_delegate`): a + three-state verdict — delegate, degrade to the standard pipeline, or raise; +- payload construction and submission (:meth:`ConnectorDelegate.submit`); +- background monitoring (:meth:`ConnectorDelegate._monitor`), closing the OV + TaskRecord from the remote task's terminal state. +""" + +from __future__ import annotations + +import asyncio +import time +from typing import Any, Awaitable, Callable, Dict, List, Optional, Set, Tuple + +import httpx + +from openviking.connector.client import ConnectorClient +from openviking.connector.routing import ( + CONNECTOR_ARGS_AUTH_CONFIG_KEY, + CONNECTOR_CREDENTIAL_ARGS, + CONNECTOR_SUPPORTED_ARGS, + credential_arg_names, + detect_connector_add_type, + is_full_commit_sha, +) +from openviking.core.content_targets import ContentTargetSpec +from openviking.resource.processing_mode import ( + DEFAULT_PROCESSING_MODE, + VECTORS_ONLY, + ProcessingMode, +) +from openviking.server.identity import RequestContext +from openviking.utils.tags import normalize_search_tags +from openviking_cli.exceptions import InternalError, InvalidArgumentError +from openviking_cli.utils import get_logger + +logger = get_logger(__name__) + + +class ConnectorDelegate: + """Routes add_resource requests to the external Connector service. + + Collaborators are injected so the delegate stays free of service + lifecycle concerns: + + - ``viking_fs``: used only for the parent-existence pre-check when + ``create_parent`` is not requested. + - ``background_tasks``: the ResourceService task set; monitors register + here so ``close_background_tasks`` keeps cancelling them on shutdown. + - ``link_reason_memory``: bound ResourceService method that links one + reason memory to the import root — shared semantics with the native + pipeline, so it is called back rather than duplicated. + """ + + def __init__( + self, + *, + viking_fs: Any, + background_tasks: Set["asyncio.Task[Any]"], + link_reason_memory: Callable[..., Awaitable[None]], + ) -> None: + self._viking_fs = viking_fs + self._background_tasks = background_tasks + self._link_reason_memory = link_reason_memory + + @staticmethod + def resolve_add_type(path: str, declared_add_type: Optional[str]) -> Optional[Tuple[str, bool]]: + """Resolve the Connector ``(add_type, connector_only)`` for *path*. + + Probes the path when no add_type is declared. A declared add_type is + validated against the probe: registry types (tos/git) must match it + exactly, and other types must not claim a path that probes as a + registry type. A declared request is an explicit routing instruction, + so it is always connector-only: it never degrades to the standard + pipeline. + """ + detected = detect_connector_add_type(path) + if declared_add_type is None: + return detected + probed = detected[0] if detected else None + if declared_add_type in CONNECTOR_SUPPORTED_ARGS: + if probed != declared_add_type: + raise InvalidArgumentError( + f"add_type '{declared_add_type}' does not match the source path" + + (f", which is a '{probed}' source" if probed else "") + + "." + ) + elif probed is not None: + raise InvalidArgumentError( + f"add_type '{declared_add_type}' does not match the source path, " + f"which is a '{probed}' source." + ) + return (declared_add_type, True) + + def should_delegate( + self, + path: str, + *, + ctx: Optional[RequestContext] = None, + declared_add_type: Optional[str] = None, + to: Optional[str] = None, + parent: Optional[str] = None, + wait: bool = False, + instruction: str = "", + build_index: bool = True, + summarize: bool = False, + processing_mode: ProcessingMode = DEFAULT_PROCESSING_MODE, + watch_interval: float = 0, + connector_args: Optional[Dict[str, Any]] = None, + kwargs: Optional[Dict[str, Any]] = None, + ) -> bool: + """Decide whether a top-level resource path belongs to Connector. + + Returns True to delegate to the Connector, False to route to the + standard pipeline. Source types are detected the same way the + standard pipeline routes them (see connector.routing), not by raw + URL scheme. A source type only the Connector can import (tos) + raises InvalidArgumentError when the Connector is disabled, does + not allow the type, or cannot honor the request parameters — + degrading such a request would only fail later with a misleading + parse error. Source types the standard pipeline can also handle + degrade to it instead when parameters are unsupported, unless the + request carries Connector-only credentials; credentialed requests + fail closed so secrets never enter a durable native job. + """ + from openviking_cli.utils.config.open_viking_config import get_openviking_config + + resolved = self.resolve_add_type(path, declared_add_type) + if resolved is None: + return False + add_type, connector_only = resolved + connector_args = connector_args or {} + credential_args = credential_arg_names(add_type, connector_args) + + config = get_openviking_config().connector + if not config.enable or add_type not in config.allowed_add_types: + if declared_add_type is not None: + raise InvalidArgumentError( + f"add_type '{add_type}' requires the Connector integration, " + "which is disabled or does not allow this source type." + ) + if credential_args: + raise InvalidArgumentError( + f"Credential args {credential_args} for '{add_type}' require Connector " + "import, but the Connector is disabled or does not allow this source type." + ) + if connector_only: + raise InvalidArgumentError( + f"'{path}' can only be imported through the Connector integration, " + "which is disabled or does not allow this source type." + ) + return False + + if add_type == "git": + from openviking.parse.accessors.git_accessor import GitAccessor + + _, _, commit = GitAccessor()._parse_repo_source(path, **connector_args) + if commit and not is_full_commit_sha(str(commit)): + # The plugin fetches pinned commits by SHA over the wire and + # cannot resolve abbreviated forms; the standard pipeline + # resolves them locally after a full clone. + if declared_add_type is not None: + raise InvalidArgumentError( + "Connector git import cannot resolve an abbreviated commit SHA; " + "pass a full 40-character SHA or drop add_type to use the " + "standard pipeline." + ) + if credential_args: + raise InvalidArgumentError( + f"Credential args {credential_args} for 'git' cannot fall back to the " + "standard import pipeline; pass a full 40-character commit SHA." + ) + logger.info( + "Connector git import degrades to the standard pipeline: " + "abbreviated commit SHA %r needs a local clone to resolve.", + commit, + ) + return False + + if ctx is not None and (to or parent): + target = ContentTargetSpec.from_fields( + ctx=ctx, + kind="resource", + to=to, + parent=parent, + create_parent=bool((kwargs or {}).get("create_parent", False)), + ) + to = target.to + parent = target.parent + + unsupported = self._unsupported_params( + add_type=add_type, + wait=wait, + instruction=instruction, + build_index=build_index, + summarize=summarize, + processing_mode=processing_mode, + watch_interval=watch_interval, + connector_args=connector_args, + kwargs=kwargs or {}, + to=to, + parent=parent, + ) + if not unsupported: + return True + detail = "; ".join(unsupported) + if connector_only: + raise InvalidArgumentError(f"Connector import does not support: {detail}") + if credential_args: + raise InvalidArgumentError( + f"Credential args {credential_args} for '{add_type}' cannot fall back to the " + f"standard import pipeline. Connector import does not support: {detail}" + ) + logger.info( + f"[ConnectorDelegate] Connector does not support {detail} for path {path}; " + "falling back to the standard import pipeline" + ) + return False + + @staticmethod + def _unsupported_params( + *, + add_type: str, + wait: bool, + instruction: str, + build_index: bool, + summarize: bool, + processing_mode: ProcessingMode, + watch_interval: float, + connector_args: Dict[str, Any], + kwargs: Dict[str, Any], + to: Optional[str] = None, + parent: Optional[str] = None, + ) -> List[str]: + """add_resource params the Connector delegation cannot honor. + + Returns an empty list when the request is fully supported. + """ + unsupported: List[str] = [] + if wait: + unsupported.append("wait=true (Connector imports run asynchronously)") + if parent: + unsupported.append("parent targets (Connector imports require an exact 'to' target)") + if not to: + unsupported.append("missing exact 'to' target") + elif to != "viking://resources" and not to.startswith("viking://resources/"): + unsupported.append("to outside the public resources root (viking://resources/...)") + if watch_interval > 0: + unsupported.append("watch_interval>0 (Connector imports cannot be watched yet)") + if instruction: + unsupported.append("instruction") + if not build_index: + unsupported.append("build_index=false") + if summarize: + unsupported.append("summarize=true") + if processing_mode == VECTORS_ONLY: + unsupported.append("processing_mode=vectors_only") + if kwargs.get("strict"): + unsupported.append("strict=true (Connector imports fail per file, not all-or-nothing)") + for field in ("ignore_dirs", "include", "exclude"): + if kwargs.get(field): + unsupported.append(f"{field} (Connector imports cannot filter the source tree)") + if kwargs.get("preserve_structure") is False: + unsupported.append( + "preserve_structure=false (Connector always mirrors the source directory tree)" + ) + if not kwargs.get("directly_upload_media", True): + unsupported.append("directly_upload_media=false") + if kwargs.get("source_name"): + unsupported.append("source_name") + if add_type in CONNECTOR_SUPPORTED_ARGS: + supported_args = CONNECTOR_SUPPORTED_ARGS.get( + add_type, frozenset() + ) | CONNECTOR_CREDENTIAL_ARGS.get(add_type, frozenset()) + unknown_args = sorted(set(connector_args) - supported_args) + if unknown_args: + supported_hint = ( + f"supported: {sorted(supported_args)}" + if supported_args + else "no args supported" + ) + unsupported.append(f"args keys {unknown_args} ({supported_hint} for '{add_type}')") + else: + # Declared add_types outside the registry forward args wholesale + # to the plugin via param_config; only the reserved credentials + # key is inspected here. + declared_auth = connector_args.get(CONNECTOR_ARGS_AUTH_CONFIG_KEY) + if declared_auth is not None and not isinstance(declared_auth, dict): + unsupported.append("args.auth_config that is not an object of credential fields") + return unsupported + + async def submit( + self, + path: str, + ctx: RequestContext, + to: Optional[str], + reason: str = "", + declared_add_type: Optional[str] = None, + connector_args: Optional[Dict[str, Any]] = None, + tags: Optional[List[str]] = None, + tag_mode: str = "replace", + **kwargs: Any, + ) -> Dict[str, Any]: + """Route add_resource to the external Connector service.""" + from openviking.service.task_tracker import get_task_tracker + from openviking_cli.utils.config.open_viking_config import get_openviking_config + + config = get_openviking_config().connector + if not ctx.api_key: + raise InvalidArgumentError("Connector import requires an API key in the request.") + + resolved = self.resolve_add_type(path, declared_add_type) + if resolved is None: + raise InvalidArgumentError(f"'{path}' does not match any Connector source type.") + add_type, _ = resolved + + target = ContentTargetSpec.from_fields( + ctx=ctx, + kind="resource", + to=to, + create_parent=bool(kwargs.get("create_parent", False)), + ) + if not target.to: + raise InvalidArgumentError("Connector import requires an exact 'to' target.") + task_resource_id = target.to + + if not kwargs.get("create_parent", False): + # Match the native pipeline default: unless create_parent=true is + # given, the parent of the exact target must pre-exist. + parent_uri, _, _ = task_resource_id.rpartition("/") + if parent_uri.startswith("viking://") and not await self._viking_fs.exists( + parent_uri, ctx + ): + raise InvalidArgumentError( + f"Parent directory '{parent_uri}' does not exist; " + "pass create_parent=true to create it." + ) + + # Credentials are separated into the dedicated auth_config field before + # forwarding: Connector and plugin logs redact auth_config while logging + # param_config verbatim. The incoming OpenViking HTTP body may still be + # captured when the explicitly unsafe observability body dump is enabled. + if add_type in CONNECTOR_SUPPORTED_ARGS: + forwarded_args = { + key: value + for key, value in (connector_args or {}).items() + if key in CONNECTOR_SUPPORTED_ARGS.get(add_type, frozenset()) and value + } + auth_config = { + key: str(value) + for key, value in (connector_args or {}).items() + if key in CONNECTOR_CREDENTIAL_ARGS.get(add_type, frozenset()) and value + } + else: + # Declared add_type outside the registry: the plugin owns the args + # semantics, so they travel wholesale through param_config. The + # caller supplies credentials under the reserved args.auth_config + # key, lifted into the top-level auth_config field. + forwarded_args = { + key: value + for key, value in (connector_args or {}).items() + if key != CONNECTOR_ARGS_AUTH_CONFIG_KEY and value is not None + } + declared_auth = (connector_args or {}).get(CONNECTOR_ARGS_AUTH_CONFIG_KEY) + if declared_auth is not None and not isinstance(declared_auth, dict): + raise InvalidArgumentError( + "args.auth_config must be an object of credential fields." + ) + auth_config = { + str(key): str(value) for key, value in (declared_auth or {}).items() if value + } + + tos_path: Optional[str] = None + param_config: Optional[Dict[str, Any]] = None + if add_type == "tos": + source_path = path[len("tos://") :].strip() + if not source_path: + raise InvalidArgumentError( + "Connector TOS import requires path='tos:///'." + ) + tos_path = source_path + elif add_type == "git": + from openviking.parse.accessors.git_accessor import GitAccessor + + # Source-specific settings stay in param_config. The exact target + # is transported separately as the top-level ``to`` field. + repo_url, branch, commit = GitAccessor()._parse_repo_source( + path.strip(), **forwarded_args + ) + param_config = {"repo_url": repo_url} + if commit: + # Routing already degraded abbreviated SHAs to the standard + # pipeline. Preserve the native accessor's historical precedence: + # when both commit and branch/ref are present, send only the full + # commit SHA so the plugin receives an unambiguous pinned request. + param_config["commit"] = str(commit) + elif branch: + param_config["branch"] = str(branch) + else: + # Declared add_type outside the registry: the plugin owns the + # source semantics. The top-level path travels under the agreed + # param_config["path"] key -- args can never carry "path", it is + # a reserved core add_resource field. + param_config = dict(forwarded_args) + param_config["path"] = path + + client = ConnectorClient( + doc_add_url=config.connector, + task_info_url=config.tracker, + account_id=ctx.account_id, + ) + extra_params = None + if tags is not None: + extra_params = { + "tags": normalize_search_tags(tags), + "tag_mode": tag_mode, + } + + task_tracker = get_task_tracker() + task = await task_tracker.create( + "connector_import", + resource_id=task_resource_id, + account_id=ctx.account_id, + user_id=ctx.user.user_id, + ) + try: + result = await client.submit_doc_add( + add_type=add_type, + api_key=ctx.api_key, + tos_path=tos_path, + to=task_resource_id, + include_child=True, + param_config=param_config, + auth_config=auth_config or None, + extra_params=extra_params, + ) + + connector_task_key = result.get("task_key") or result.get("TaskKey") or "" + if not connector_task_key: + raise InternalError( + f"Connector accepted the import but returned no task key: {result}" + ) + except asyncio.CancelledError: + await task_tracker.fail( + task.task_id, + "connector task submission cancelled", + account_id=ctx.account_id, + user_id=ctx.user.user_id, + ) + raise + except Exception as exc: + await task_tracker.fail( + task.task_id, + str(exc), + account_id=ctx.account_id, + user_id=ctx.user.user_id, + ) + raise + + monitor = self._monitor( + client=client, + connector_task_key=connector_task_key, + ov_task_id=task.task_id, + poll_interval_ms=config.poll_interval_ms, + timeout_seconds=config.timeout_seconds, + ctx=ctx, + reason=reason, + link_root_uri=task_resource_id or "viking://resources", + ) + + background = asyncio.create_task(monitor) + self._background_tasks.add(background) + background.add_done_callback(self._background_tasks.discard) + + response = { + "status": "accepted", + "task_id": task.task_id, + "connector_task_key": connector_task_key, + } + if task_resource_id: + response["resource_id"] = task_resource_id + return response + + async def _monitor( + self, + client: ConnectorClient, + connector_task_key: str, + ov_task_id: str, + poll_interval_ms: int, + timeout_seconds: int, + ctx: RequestContext, + reason: str = "", + link_root_uri: str = "", + ) -> Dict[str, Any]: + """Poll the Connector task until terminal state, then update OV TaskRecord. + + Returns a terminal payload: ``{"status": "completed", ...}`` on + success, ``{"status": "failed", "error": ...}`` otherwise. When + *reason* is set, a successful import links one reason memory to + *link_root_uri* (the import root), matching the native add_resource + reason semantics of one reason entry referencing the resource root. + """ + from openviking.service.task_tracker import get_task_tracker + + task_tracker = get_task_tracker() + await task_tracker.start( + ov_task_id, + account_id=ctx.account_id, + user_id=ctx.user.user_id, + ) + + poll_interval = poll_interval_ms / 1000.0 + deadline = time.perf_counter() + timeout_seconds + terminal_statuses = {"succeeded", "failed", "cancelled"} + + try: + while time.perf_counter() < deadline: + await asyncio.sleep(poll_interval) + try: + info = await client.get_task_info(connector_task_key, ctx.api_key) + except httpx.HTTPStatusError as exc: + status_code = exc.response.status_code + if status_code not in {408, 429} and status_code < 500: + raise + logger.warning( + "[ConnectorDelegate] Transient Connector task polling HTTP error " + f"for {connector_task_key}: {status_code}; retrying" + ) + continue + except httpx.RequestError as exc: + logger.warning( + "[ConnectorDelegate] Transient Connector task polling error " + f"for {connector_task_key}: {exc}; retrying" + ) + continue + status = (info.get("Status") or info.get("status") or "").lower() + + await task_tracker.update_stage( + ov_task_id, + f"connector:{status}", + account_id=ctx.account_id, + user_id=ctx.user.user_id, + ) + + if status in terminal_statuses: + if status == "succeeded": + completion: Dict[str, Any] = { + "connector_status": status, + "connector_task_key": connector_task_key, + } + if (reason or "").strip() and link_root_uri: + link_result: Dict[str, Any] = {"root_uri": link_root_uri} + await self._link_reason_memory( + result=link_result, + ctx=ctx, + reason=reason, + source_name=None, + timeout=None, + ) + for key in ("memory_linking", "warnings"): + if key in link_result: + completion[key] = link_result[key] + await task_tracker.complete( + ov_task_id, + completion, + account_id=ctx.account_id, + user_id=ctx.user.user_id, + ) + return {"status": "completed", **completion} + error_msg = info.get("ErrorMessage") or info.get("error_message") or status + failure = f"connector task {status}: {error_msg}" + await task_tracker.fail( + ov_task_id, + failure, + account_id=ctx.account_id, + user_id=ctx.user.user_id, + ) + return {"status": "failed", "error": failure} + + timeout_msg = f"connector task timed out after {timeout_seconds}s" + await task_tracker.fail( + ov_task_id, + timeout_msg, + account_id=ctx.account_id, + user_id=ctx.user.user_id, + ) + return {"status": "failed", "error": timeout_msg} + except asyncio.CancelledError: + await task_tracker.fail( + ov_task_id, + "background connector task monitoring cancelled", + account_id=ctx.account_id, + user_id=ctx.user.user_id, + ) + raise + except Exception as exc: + logger.error(f"[ConnectorDelegate] Connector task monitor error: {exc}") + await task_tracker.fail( + ov_task_id, + str(exc), + account_id=ctx.account_id, + user_id=ctx.user.user_id, + ) + return {"status": "failed", "error": str(exc)} diff --git a/openviking/connector/routing.py b/openviking/connector/routing.py index 347b84d02..7ff260640 100644 --- a/openviking/connector/routing.py +++ b/openviking/connector/routing.py @@ -12,7 +12,7 @@ hosting ``http(s)`` URLs). from __future__ import annotations -from typing import Dict, FrozenSet, Optional, Tuple +from typing import Dict, FrozenSet, List, Optional, Tuple from openviking.utils import is_git_repo_url @@ -30,9 +30,10 @@ CONNECTOR_SUPPORTED_ARGS: Dict[str, FrozenSet[str]] = { # add_resource ``args`` keys carrying source credentials (git: HTTPS PAT as # the Basic-Auth password, username defaults to "oauth2" plugin-side). They # are stripped out of args and transported in the top-level ``auth_config`` -# request field -- never ``param_config`` -- because only the auth channel -# is excluded from request logging on every hop (Connector and plugin log -# param_config verbatim). One-shot use: credentials are never persisted, and +# request field -- never ``param_config`` -- so Connector and plugin logs +# redact them while logging param_config verbatim. The incoming OpenViking +# HTTP body may still be captured when the explicitly unsafe observability +# body dump is enabled. One-shot use: credentials are never persisted, and # requests carrying them must not fall back to a durable native import job. CONNECTOR_CREDENTIAL_ARGS: Dict[str, FrozenSet[str]] = { "tos": frozenset(), @@ -40,6 +41,22 @@ CONNECTOR_CREDENTIAL_ARGS: Dict[str, FrozenSet[str]] = { } +# Reserved ``args`` key for source credentials of declared add_types outside +# the registries above (e.g. add_type="feishu"). OpenViking cannot know each +# plugin's credential fields, so the caller supplies them as a mapping under +# this key; it is lifted verbatim into the top-level ``auth_config`` request +# field and never merged into param_config. All other args keys travel to the +# plugin wholesale through param_config. +CONNECTOR_ARGS_AUTH_CONFIG_KEY = "auth_config" + + +def credential_arg_names(add_type: str, args: Optional[Dict[str, object]]) -> List[str]: + """Sorted credential-arg keys of *add_type* that are present in *args*.""" + return sorted( + set(args or ()).intersection(CONNECTOR_CREDENTIAL_ARGS.get(add_type, frozenset())) + ) + + def is_full_commit_sha(ref: str) -> bool: """True when *ref* is a full 40-hex commit SHA. diff --git a/openviking/server/local_input_guard.py b/openviking/server/local_input_guard.py index 559c832c9..de2009523 100644 --- a/openviking/server/local_input_guard.py +++ b/openviking/server/local_input_guard.py @@ -55,9 +55,17 @@ def looks_like_local_path(value: str) -> bool: ) -def require_remote_resource_source(source: str) -> str: - """Reject direct host-path resource ingestion over HTTP.""" - if not is_remote_resource_source(source): +def require_remote_resource_source( + source: str, *, declared_connector_add_type: Optional[str] = None +) -> str: + """Reject direct host-path resource ingestion over HTTP. + + A declared Connector add_type skips the remote-shape requirement: such a + request is either delegated to the Connector or rejected by the routing + predicate with a clear error, and never enters local path resolution. + URL-shaped sources still must point at public remote hosts. + """ + if declared_connector_add_type is None and not is_remote_resource_source(source): raise PermissionDeniedError( "HTTP server only accepts remote resource URLs or temp-uploaded files; " "direct host filesystem paths are not allowed." diff --git a/openviking/server/mcp_endpoint.py b/openviking/server/mcp_endpoint.py index 147bde061..7f8539e83 100644 --- a/openviking/server/mcp_endpoint.py +++ b/openviking/server/mcp_endpoint.py @@ -548,6 +548,7 @@ async def _maybe_sitemap_hint(path: str) -> str: async def add_resource( path: str = "", temp_file_id: str = "", + add_type: str = "", description: str = "", watch_interval: float = 0, processing_mode: ProcessingMode = DEFAULT_PROCESSING_MODE, @@ -570,6 +571,12 @@ async def add_resource( Args: path: Remote URL or local filesystem path. Required unless ``temp_file_id`` is set. temp_file_id: Server-minted upload id from a prior signed local-file upload. + add_type: Explicit Connector source type (e.g. "tos", "git"). When set, the + request routes through the Connector integration (must be enabled + server-side) and ``path`` is sent verbatim — never treated as a local + file. Requires an exact ``to`` target and cannot be combined with + ``temp_file_id`` or ``parent``. Leave empty for the default path-probing + behavior. description: Optional human-readable reason for adding the resource. watch_interval: Auto-refresh cadence in minutes. 0 = no watch. Prefer >=1440 (24h) unless the source changes faster — every refresh re-embeds the whole resource. @@ -577,9 +584,10 @@ async def add_resource( processing_mode: "semantic_and_vectors" for normal semantic processing, or "vectors_only" to skip semantic understanding and only build vector indexes. to: Target URI under viking://resources/ (e.g. "viking://resources/volcengine/OpenViking"). - Leave empty to derive a URI from the source. + Required when ``add_type`` is set; otherwise leave empty to derive a URI + from the source. parent: Parent URI under viking://resources/ for remote imports. Mutually exclusive - with ``to``. + with ``to`` and not supported when ``add_type`` is set. tags: Optional explicit k=v retrieval tags to apply after ingestion. tag_mode: Tag update mode, "replace" or "append". Defaults to "replace". args: Parser-specific options, e.g. {"feishu_access_token": "..."} for Feishu imports, @@ -596,6 +604,16 @@ async def add_resource( "use a positive number of minutes (>=1440 recommended) to subscribe to auto-refresh." ) + add_type = add_type.strip() + if add_type and temp_file_id: + return "Error: add_type cannot be combined with temp_file_id." + if add_type and not path: + return "Error: add_type requires 'path'." + if add_type and parent: + return "Error: add_type cannot be combined with parent." + if add_type and not to: + return "Error: add_type requires an exact 'to' target." + # Branch 1: ingest by temp_file_id. Kept for backward compat / REST-style use — the # signed upload now auto-ingests server-side, so agents no longer need this second leg. if temp_file_id: @@ -646,13 +664,18 @@ async def add_resource( f'Pass it as the temp_file_id kwarg: add_resource(temp_file_id="{path}")' ) - # Branch 3: remote URL — same flow as before - if is_remote_resource_source(path): + # Branch 3: remote URL, or an explicitly declared Connector source type + # (declared requests are delegated or rejected server-side, never resolved + # as local paths) + if add_type or is_remote_resource_source(path): try: - path = require_remote_resource_source(path) + path = require_remote_resource_source( + path, declared_connector_add_type=add_type or None + ) result = await service.resources.add_resource( path=path, ctx=ctx, + add_type=add_type or None, to=to or None, parent=parent or None, reason=description, @@ -683,7 +706,8 @@ async def add_resource( # Detect-and-suggest: if this single page belongs to a site that exposes a # sitemap/RSS feed, hint at whole-site ingestion. Never auto-crawls; the # add above is already done, so a slow/failed probe has no functional impact. - hint = await _maybe_sitemap_hint(path) + # Declared Connector imports ingest the whole source already — no hint. + hint = None if add_type else await _maybe_sitemap_hint(path) if hint: message += "\n" + hint return message diff --git a/openviking/server/routers/resources.py b/openviking/server/routers/resources.py index 4376e113c..eaa142aa3 100644 --- a/openviking/server/routers/resources.py +++ b/openviking/server/routers/resources.py @@ -31,10 +31,17 @@ class AddResourceRequest(BaseModel): Either path or temp_file_id must be provided. temp_file_id: Temporary upload id returned by /api/v1/resources/temp_upload. Either path or temp_file_id must be provided. + add_type: Explicit Connector source type (e.g. "tos", "git"). When set, the + request routes to the Connector integration: the type must be enabled in + connector.allowed_add_types, and args are forwarded to the source plugin + (credentials under args.auth_config). Never degrades to the standard + pipeline. Requires 'path' and an exact 'to' target; cannot be combined + with 'temp_file_id' or 'parent'. to: Target URI for the resource (e.g., "viking://resources/my_resource"). - If not specified, an auto-generated URI will be used. + Required when add_type is set. Otherwise, if not specified, an + auto-generated URI will be used. parent: Parent URI under which the resource will be stored. - Cannot be used together with 'to'. + Cannot be used together with 'to' or 'add_type'. create_parent: Whether to automatically create the parent directory if it doesn't exist. Default is False. reason: Reason for adding the resource. Used for documentation and monitoring. @@ -69,6 +76,7 @@ class AddResourceRequest(BaseModel): path: Optional[str] = None temp_file_id: Optional[str] = None + add_type: Optional[str] = None to: Optional[str] = None parent: Optional[str] = None create_parent: bool = False @@ -96,6 +104,20 @@ class AddResourceRequest(BaseModel): raise ValueError("Either 'path' or 'temp_file_id' must be provided") return self + @model_validator(mode="after") + def check_add_type(self): + if self.add_type is not None: + self.add_type = self.add_type.strip() or None + if self.add_type and self.temp_file_id: + raise ValueError("'add_type' cannot be combined with 'temp_file_id'") + if self.add_type and not self.path: + raise ValueError("'add_type' requires 'path'") + if self.add_type and self.parent: + raise ValueError("'add_type' cannot be combined with 'parent'") + if self.add_type and not self.to: + raise ValueError("'add_type' requires an exact 'to' target") + return self + class AddSkillRequest(BaseModel): """Request model for add_skill. @@ -209,7 +231,7 @@ async def add_resource( original_filename = resolved.original_filename allow_local_path_resolution = True elif path is not None: - path = require_remote_resource_source(path) + path = require_remote_resource_source(path, declared_connector_add_type=request.add_type) if path is None: raise InvalidArgumentError("Either 'path' or 'temp_file_id' must be provided.") @@ -243,6 +265,7 @@ async def add_resource( result = await service.resources.add_resource( path=path, ctx=_ctx, + add_type=request.add_type, to=request.to, parent=request.parent, reason=request.reason, diff --git a/openviking/service/resource_service.py b/openviking/service/resource_service.py index c9676d9d1..3036d0a91 100644 --- a/openviking/service/resource_service.py +++ b/openviking/service/resource_service.py @@ -17,11 +17,7 @@ from pathlib import Path from typing import TYPE_CHECKING, Any, Dict, List, Optional from uuid import uuid4 -import httpx - -from openviking.connector.client import ConnectorClient from openviking.core.content_targets import ContentTargetSpec -from openviking.utils.ingest_options import IngestOptions from openviking.core.uri_validation import validate_optional_content_target_uri from openviking.parse.parsers.constants import MPEG_TS_EXTENSION_ALIAS from openviking.resource.feishu_watch_auth import ( @@ -32,7 +28,6 @@ from openviking.resource.feishu_watch_auth import ( ) from openviking.resource.processing_mode import ( DEFAULT_PROCESSING_MODE, - VECTORS_ONLY, ProcessingMode, normalize_processing_mode, ) @@ -58,6 +53,7 @@ from openviking.telemetry.resource_summary import ( unregister_wait_telemetry, ) from openviking.utils import is_git_repo_url, parse_code_hosting_url +from openviking.utils.ingest_options import IngestOptions from openviking.utils.media_processor import _smart_stem from openviking.utils.network_guard import ensure_public_remote_target from openviking.utils.resource_processor import ResourceProcessor @@ -72,6 +68,7 @@ from openviking_cli.exceptions import ( from openviking_cli.utils import get_logger if TYPE_CHECKING: + from openviking.connector.delegate import ConnectorDelegate from openviking.parse.accessors.base import LocalResource from openviking.resource.watch_manager import WatchManager from openviking.resource.watch_scheduler import WatchScheduler @@ -168,6 +165,7 @@ class ResourceService: self._watch_scheduler = watch_scheduler self._resource_memory_link_service = resource_memory_link_service self._background_tasks: set[asyncio.Task[Any]] = set() + self._connector_delegate: Optional["ConnectorDelegate"] = None def set_dependencies( self, @@ -577,11 +575,9 @@ class ResourceService: self._validate_add_resource_tag_policy(tags=tags, tag_mode=tag_mode) normalized_args = self._normalize_add_resource_args(args, watch_interval=watch_interval) kwargs.update(normalized_args.processor_kwargs) - from openviking.connector.routing import CONNECTOR_CREDENTIAL_ARGS + from openviking.connector.routing import credential_arg_names - credential_args = sorted( - set(kwargs).intersection(CONNECTOR_CREDENTIAL_ARGS.get("git", frozenset())) - ) + credential_args = credential_arg_names("git", kwargs) if credential_args: raise InvalidArgumentError( "Git credential args " @@ -780,6 +776,7 @@ class ResourceService: tag_mode: str = "replace", allow_local_path_resolution: bool = True, enforce_public_remote_targets: bool = False, + add_type: Optional[str] = None, args: Optional[Dict[str, Any]] = None, **kwargs, ) -> Dict[str, Any]: @@ -790,9 +787,12 @@ class ResourceService: "add_resource does not accept internal execution fields: " + ", ".join(internal_fields) ) + if isinstance(add_type, str): + add_type = add_type.strip() or None return await self._submit_resource_ingestion( path=path, ctx=ctx, + add_type=add_type, to=to, parent=parent, reason=reason, @@ -854,6 +854,7 @@ class ResourceService: self, path: str, ctx: RequestContext, + add_type: Optional[str] = None, to: Optional[str] = None, parent: Optional[str] = None, reason: str = "", @@ -876,8 +877,15 @@ class ResourceService: Args: path: Resource path (local file or URL) - to: Target URI (e.g., "viking://resources/my_resource") - parent: Parent URI under which the resource will be stored + add_type: Explicitly declared Connector source type. Routes the + request to the Connector integration without probing the path; + the type must be enabled in connector.allowed_add_types. A + declared request never degrades to the standard pipeline and + requires an exact ``to`` target. + to: Target URI (e.g., "viking://resources/my_resource"). Required + when ``add_type`` is set. + parent: Parent URI under which the resource will be stored. Not + supported when ``add_type`` is set. reason: Reason for adding the resource instruction: Processing instruction for semantic extraction wait: Whether to wait for semantic extraction and vectorization to complete @@ -937,9 +945,10 @@ class ResourceService: parent = default_parent kwargs["create_parent"] = True - if self._should_use_connector( + if self._connector.should_delegate( path, ctx=ctx, + declared_add_type=add_type, to=to, parent=parent, wait=wait, @@ -950,14 +959,16 @@ class ResourceService: watch_interval=watch_interval, connector_args=args or {}, kwargs=kwargs, - tags=tags, ): - return await self._add_resource_via_connector( + return await self._connector.submit( path=path, ctx=ctx, + declared_add_type=add_type, to=to, reason=reason, connector_args=args or {}, + tags=tags, + tag_mode=tag_mode, **kwargs, ) @@ -1495,482 +1506,18 @@ class ResourceService: # ── Connector routing ── - def _should_use_connector( - self, - path: str, - *, - ctx: Optional[RequestContext] = None, - to: Optional[str] = None, - parent: Optional[str] = None, - wait: bool = False, - instruction: str = "", - build_index: bool = True, - summarize: bool = False, - processing_mode: ProcessingMode = DEFAULT_PROCESSING_MODE, - watch_interval: float = 0, - connector_args: Optional[Dict[str, Any]] = None, - kwargs: Optional[Dict[str, Any]] = None, - tags: Optional[List[str]] = None, - ) -> bool: - """Decide whether a top-level resource path belongs to Connector. + @property + def _connector(self) -> "ConnectorDelegate": + """Connector delegation (lazy: viking_fs may be injected after init).""" + if self._connector_delegate is None: + from openviking.connector.delegate import ConnectorDelegate - Returns True to delegate to the Connector, False to route to the - standard pipeline. Source types are detected the same way the - standard pipeline routes them (see connector.routing), not by raw - URL scheme. A source type only the Connector can import (tos) - raises InvalidArgumentError when the Connector is disabled, does - not allow the type, or cannot honor the request parameters — - degrading such a request would only fail later with a misleading - parse error. Source types the standard pipeline can also handle - degrade to it instead when parameters are unsupported, unless the - request carries Connector-only credentials; credentialed requests - fail closed so secrets never enter a durable native job. - """ - from openviking.connector.routing import ( - CONNECTOR_CREDENTIAL_ARGS, - detect_connector_add_type, - is_full_commit_sha, - ) - from openviking_cli.utils.config.open_viking_config import get_openviking_config - - detected = detect_connector_add_type(path) - if detected is None: - return False - add_type, connector_only = detected - connector_args = connector_args or {} - credential_args = sorted( - set(connector_args).intersection(CONNECTOR_CREDENTIAL_ARGS.get(add_type, frozenset())) - ) - - config = get_openviking_config().connector - if not config.enable or add_type not in config.allowed_add_types: - if credential_args: - raise InvalidArgumentError( - f"Credential args {credential_args} for '{add_type}' require Connector " - "import, but the Connector is disabled or does not allow this source type." - ) - if connector_only: - raise InvalidArgumentError( - f"'{path}' can only be imported through the Connector integration, " - "which is disabled or does not allow this source type." - ) - return False - - if add_type == "git": - from openviking.parse.accessors.git_accessor import GitAccessor - - _, _, commit = GitAccessor()._parse_repo_source(path, **connector_args) - if commit and not is_full_commit_sha(str(commit)): - # The plugin fetches pinned commits by SHA over the wire and - # cannot resolve abbreviated forms; the standard pipeline - # resolves them locally after a full clone. - if credential_args: - raise InvalidArgumentError( - f"Credential args {credential_args} for 'git' cannot fall back to the " - "standard import pipeline; pass a full 40-character commit SHA." - ) - logger.info( - "Connector git import degrades to the standard pipeline: " - "abbreviated commit SHA %r needs a local clone to resolve.", - commit, - ) - return False - - if ctx is not None and (to or parent): - target = ContentTargetSpec.from_fields( - ctx=ctx, - kind="resource", - to=to, - parent=parent, - create_parent=bool((kwargs or {}).get("create_parent", False)), + self._connector_delegate = ConnectorDelegate( + viking_fs=self._viking_fs, + background_tasks=self._background_tasks, + link_reason_memory=self._link_resource_reason_memory, ) - to = target.to - parent = target.parent - - unsupported = self._unsupported_connector_params( - add_type=add_type, - wait=wait, - instruction=instruction, - build_index=build_index, - summarize=summarize, - processing_mode=processing_mode, - watch_interval=watch_interval, - connector_args=connector_args, - kwargs=kwargs or {}, - to=to, - parent=parent, - tags=tags, - ) - if not unsupported: - return True - detail = "; ".join(unsupported) - if connector_only: - raise InvalidArgumentError(f"Connector import does not support: {detail}") - if credential_args: - raise InvalidArgumentError( - f"Credential args {credential_args} for '{add_type}' cannot fall back to the " - f"standard import pipeline. Connector import does not support: {detail}" - ) - logger.info( - f"[ResourceService] Connector does not support {detail} for path {path}; " - "falling back to the standard import pipeline" - ) - return False - - @staticmethod - def _unsupported_connector_params( - *, - add_type: str, - wait: bool, - instruction: str, - build_index: bool, - summarize: bool, - processing_mode: ProcessingMode, - watch_interval: float, - connector_args: Dict[str, Any], - kwargs: Dict[str, Any], - to: Optional[str] = None, - parent: Optional[str] = None, - tags: Optional[List[str]] = None, - ) -> List[str]: - """add_resource params the Connector delegation cannot honor. - - Returns an empty list when the request is fully supported. - """ - from openviking.connector.routing import ( - CONNECTOR_CREDENTIAL_ARGS, - CONNECTOR_SUPPORTED_ARGS, - ) - - unsupported: List[str] = [] - if wait: - unsupported.append("wait=true (Connector imports run asynchronously)") - if parent: - unsupported.append("parent targets (Connector imports require an exact 'to' target)") - if not to: - unsupported.append("missing exact 'to' target") - elif to != "viking://resources" and not to.startswith("viking://resources/"): - unsupported.append("to outside the public resources root (viking://resources/...)") - if watch_interval > 0: - unsupported.append("watch_interval>0 (Connector imports cannot be watched yet)") - if instruction: - unsupported.append("instruction") - if not build_index: - unsupported.append("build_index=false") - if summarize: - unsupported.append("summarize=true") - if processing_mode == VECTORS_ONLY: - unsupported.append("processing_mode=vectors_only") - if kwargs.get("strict"): - unsupported.append("strict=true (Connector imports fail per file, not all-or-nothing)") - for field in ("ignore_dirs", "include", "exclude"): - if kwargs.get(field): - unsupported.append(f"{field} (Connector imports cannot filter the source tree)") - if kwargs.get("preserve_structure") is False: - unsupported.append( - "preserve_structure=false (Connector always mirrors the source directory tree)" - ) - if not kwargs.get("directly_upload_media", True): - unsupported.append("directly_upload_media=false") - if kwargs.get("source_name"): - unsupported.append("source_name") - if tags is not None: - unsupported.append("tags (Connector imports cannot apply ingestion tags yet)") - supported_args = CONNECTOR_SUPPORTED_ARGS.get( - add_type, frozenset() - ) | CONNECTOR_CREDENTIAL_ARGS.get(add_type, frozenset()) - unknown_args = sorted(set(connector_args) - supported_args) - if unknown_args: - supported_hint = ( - f"supported: {sorted(supported_args)}" if supported_args else "no args supported" - ) - unsupported.append(f"args keys {unknown_args} ({supported_hint} for '{add_type}')") - return unsupported - - async def _add_resource_via_connector( - self, - path: str, - ctx: RequestContext, - to: Optional[str], - reason: str = "", - connector_args: Optional[Dict[str, Any]] = None, - **kwargs: Any, - ) -> Dict[str, Any]: - """Route add_resource to the external Connector service.""" - from openviking.connector.routing import ( - CONNECTOR_CREDENTIAL_ARGS, - CONNECTOR_SUPPORTED_ARGS, - detect_connector_add_type, - ) - from openviking.service.task_tracker import get_task_tracker - from openviking_cli.utils.config.open_viking_config import get_openviking_config - - config = get_openviking_config().connector - if not ctx.api_key: - raise InvalidArgumentError("Connector import requires an API key in the request.") - - detected = detect_connector_add_type(path) - if detected is None: - raise InvalidArgumentError(f"'{path}' does not match any Connector source type.") - add_type, _ = detected - - target = ContentTargetSpec.from_fields( - ctx=ctx, - kind="resource", - to=to, - create_parent=bool(kwargs.get("create_parent", False)), - ) - if not target.to: - raise InvalidArgumentError("Connector import requires an exact 'to' target.") - task_resource_id = target.to - - if not kwargs.get("create_parent", False): - # Match the native pipeline default: unless create_parent=true is - # given, the parent of the exact target must pre-exist. - parent_uri, _, _ = task_resource_id.rpartition("/") - if parent_uri.startswith("viking://") and not await self._viking_fs.exists( - parent_uri, ctx - ): - raise InvalidArgumentError( - f"Parent directory '{parent_uri}' does not exist; " - "pass create_parent=true to create it." - ) - - forwarded_args = { - key: value - for key, value in (connector_args or {}).items() - if key in CONNECTOR_SUPPORTED_ARGS.get(add_type, frozenset()) and value - } - # Credentials travel exclusively in the dedicated auth_config field: - # param_config is logged verbatim by the Connector and the plugin, - # auth_config is redacted on every hop. - auth_config = { - key: str(value) - for key, value in (connector_args or {}).items() - if key in CONNECTOR_CREDENTIAL_ARGS.get(add_type, frozenset()) and value - } - - tos_path: Optional[str] = None - param_config: Optional[Dict[str, Any]] = None - if add_type == "tos": - source_path = path[len("tos://") :].strip() - if not source_path: - raise InvalidArgumentError( - "Connector TOS import requires path='tos:///'." - ) - tos_path = source_path - elif add_type == "git": - from openviking.parse.accessors.git_accessor import GitAccessor - - # Source-specific settings stay in param_config. The exact target - # is transported separately as the top-level ``to`` field. - repo_url, branch, commit = GitAccessor()._parse_repo_source( - path.strip(), **forwarded_args - ) - param_config = {"repo_url": repo_url} - if commit: - # Routing already degraded abbreviated SHAs to the standard - # pipeline. Preserve the native accessor's historical precedence: - # when both commit and branch/ref are present, send only the full - # commit SHA so the plugin receives an unambiguous pinned request. - param_config["commit"] = str(commit) - elif branch: - param_config["branch"] = str(branch) - else: # pragma: no cover - detection only yields tos/git today - raise InvalidArgumentError(f"Connector add_type '{add_type}' is not supported.") - - client = ConnectorClient( - doc_add_url=config.connector, - task_info_url=config.tracker, - account_id=ctx.account_id, - ) - - task_tracker = get_task_tracker() - task = await task_tracker.create( - "connector_import", - resource_id=task_resource_id, - account_id=ctx.account_id, - user_id=ctx.user.user_id, - ) - try: - result = await client.submit_doc_add( - add_type=add_type, - api_key=ctx.api_key, - tos_path=tos_path, - to=task_resource_id, - include_child=True, - param_config=param_config, - auth_config=auth_config or None, - extra_params=None, - ) - - connector_task_key = result.get("task_key") or result.get("TaskKey") or "" - if not connector_task_key: - raise InternalError( - f"Connector accepted the import but returned no task key: {result}" - ) - except asyncio.CancelledError: - await task_tracker.fail( - task.task_id, - "connector task submission cancelled", - account_id=ctx.account_id, - user_id=ctx.user.user_id, - ) - raise - except Exception as exc: - await task_tracker.fail( - task.task_id, - str(exc), - account_id=ctx.account_id, - user_id=ctx.user.user_id, - ) - raise - - monitor = self._monitor_connector_task( - client=client, - connector_task_key=connector_task_key, - ov_task_id=task.task_id, - poll_interval_ms=config.poll_interval_ms, - timeout_seconds=config.timeout_seconds, - ctx=ctx, - reason=reason, - link_root_uri=task_resource_id or "viking://resources", - ) - - background = asyncio.create_task(monitor) - self._background_tasks.add(background) - background.add_done_callback(self._background_tasks.discard) - - response = { - "status": "accepted", - "task_id": task.task_id, - "connector_task_key": connector_task_key, - } - if task_resource_id: - response["resource_id"] = task_resource_id - return response - - async def _monitor_connector_task( - self, - client: ConnectorClient, - connector_task_key: str, - ov_task_id: str, - poll_interval_ms: int, - timeout_seconds: int, - ctx: RequestContext, - reason: str = "", - link_root_uri: str = "", - ) -> Dict[str, Any]: - """Poll the Connector task until terminal state, then update OV TaskRecord. - - Returns a terminal payload: ``{"status": "completed", ...}`` on - success, ``{"status": "failed", "error": ...}`` otherwise. When - *reason* is set, a successful import links one reason memory to - *link_root_uri* (the import root), matching the native add_resource - reason semantics of one reason entry referencing the resource root. - """ - from openviking.service.task_tracker import get_task_tracker - - task_tracker = get_task_tracker() - await task_tracker.start( - ov_task_id, - account_id=ctx.account_id, - user_id=ctx.user.user_id, - ) - - poll_interval = poll_interval_ms / 1000.0 - deadline = time.perf_counter() + timeout_seconds - terminal_statuses = {"succeeded", "failed", "cancelled"} - - try: - while time.perf_counter() < deadline: - await asyncio.sleep(poll_interval) - try: - info = await client.get_task_info(connector_task_key, ctx.api_key) - except httpx.HTTPStatusError as exc: - status_code = exc.response.status_code - if status_code not in {408, 429} and status_code < 500: - raise - logger.warning( - "[ResourceService] Transient Connector task polling HTTP error " - f"for {connector_task_key}: {status_code}; retrying" - ) - continue - except httpx.RequestError as exc: - logger.warning( - "[ResourceService] Transient Connector task polling error " - f"for {connector_task_key}: {exc}; retrying" - ) - continue - status = (info.get("Status") or info.get("status") or "").lower() - - await task_tracker.update_stage( - ov_task_id, - f"connector:{status}", - account_id=ctx.account_id, - user_id=ctx.user.user_id, - ) - - if status in terminal_statuses: - if status == "succeeded": - completion: Dict[str, Any] = { - "connector_status": status, - "connector_task_key": connector_task_key, - } - if (reason or "").strip() and link_root_uri: - link_result: Dict[str, Any] = {"root_uri": link_root_uri} - await self._link_resource_reason_memory( - result=link_result, - ctx=ctx, - reason=reason, - source_name=None, - timeout=None, - ) - for key in ("memory_linking", "warnings"): - if key in link_result: - completion[key] = link_result[key] - await task_tracker.complete( - ov_task_id, - completion, - account_id=ctx.account_id, - user_id=ctx.user.user_id, - ) - return {"status": "completed", **completion} - error_msg = info.get("ErrorMessage") or info.get("error_message") or status - failure = f"connector task {status}: {error_msg}" - await task_tracker.fail( - ov_task_id, - failure, - account_id=ctx.account_id, - user_id=ctx.user.user_id, - ) - return {"status": "failed", "error": failure} - - timeout_msg = f"connector task timed out after {timeout_seconds}s" - await task_tracker.fail( - ov_task_id, - timeout_msg, - account_id=ctx.account_id, - user_id=ctx.user.user_id, - ) - return {"status": "failed", "error": timeout_msg} - except asyncio.CancelledError: - await task_tracker.fail( - ov_task_id, - "background connector task monitoring cancelled", - account_id=ctx.account_id, - user_id=ctx.user.user_id, - ) - raise - except Exception as exc: - logger.error(f"[ResourceService] Connector task monitor error: {exc}") - await task_tracker.fail( - ov_task_id, - str(exc), - account_id=ctx.account_id, - user_id=ctx.user.user_id, - ) - return {"status": "failed", "error": str(exc)} + return self._connector_delegate async def _handle_watch_task_creation( self, diff --git a/openviking/sync_client.py b/openviking/sync_client.py index f955fc89f..24b3f3acc 100644 --- a/openviking/sync_client.py +++ b/openviking/sync_client.py @@ -238,6 +238,7 @@ class SyncOpenViking: args: Optional[Dict[str, Any]] = None, telemetry: TelemetryRequest = False, processing_mode: str = "semantic_and_vectors", + add_type: Optional[str] = None, tags: Optional[List[str]] = None, tag_mode: str = "replace", **kwargs, @@ -250,6 +251,9 @@ class SyncOpenViking: refreshed. Args: + add_type: Explicit Connector source type. Requires an exact ``to`` + target and cannot be combined with ``parent``. The source + ``path`` is forwarded verbatim. to: Exact target URI. Existing targets keep the add_resource incremental-update behavior. parent: Target parent URI for automatic child naming. build_index: Whether to build vector index immediately (default: True). @@ -257,11 +261,18 @@ class SyncOpenViking: **kwargs: Extra options forwarded to the parser chain, e.g. ``strict``, ``ignore_dirs``, ``include``, ``exclude``. """ + if add_type is not None: + add_type = add_type.strip() or None + if add_type and parent: + raise ValueError("'add_type' cannot be combined with 'parent'.") + if add_type and not to: + raise ValueError("'add_type' requires an exact 'to' target.") if to and parent: raise ValueError("Cannot specify both 'to' and 'parent' at the same time.") return run_async( self._async_client.add_resource( path=path, + add_type=add_type, to=to, parent=parent, reason=reason, diff --git a/openviking_cli/client/base.py b/openviking_cli/client/base.py index f016083cd..e2ed1ae86 100644 --- a/openviking_cli/client/base.py +++ b/openviking_cli/client/base.py @@ -46,10 +46,15 @@ class BaseClient(ABC): processing_mode: str = "semantic_and_vectors", args: Optional[Dict[str, Any]] = None, telemetry: TelemetryRequest = False, + add_type: Optional[str] = None, tags: Optional[List[str]] = None, tag_mode: str = "replace", ) -> Dict[str, Any]: - """Add resource to OpenViking.""" + """Add resource to OpenViking. + + ``add_type`` declares a Connector source and requires an exact ``to`` + target; it cannot be combined with ``parent``. + """ ... @abstractmethod diff --git a/sdk/python/openviking_sdk/client.py b/sdk/python/openviking_sdk/client.py index 82f49b4e7..40907d8ab 100644 --- a/sdk/python/openviking_sdk/client.py +++ b/sdk/python/openviking_sdk/client.py @@ -642,13 +642,21 @@ class AsyncHTTPClient: args: Optional[Dict[str, Any]] = None, telemetry: Any = False, processing_mode: Optional[str] = None, + add_type: Optional[str] = None, tags: Optional[List[str]] = None, tag_mode: str = "replace", ) -> Dict[str, Any]: + if add_type is not None: + add_type = add_type.strip() or None + if add_type and parent: + raise ValueError("'add_type' cannot be combined with 'parent'.") + if add_type and not to: + raise ValueError("'add_type' requires an exact 'to' target.") if to and parent: raise ValueError("Cannot specify both 'to' and 'parent' at the same time.") request_data = { + "add_type": add_type, "to": to, "parent": parent, "reason": reason, @@ -673,7 +681,7 @@ class AsyncHTTPClient: request_data["preserve_structure"] = preserve_structure path_obj = Path(path) - if path_obj.exists(): + if not add_type and path_obj.exists(): if path_obj.is_dir(): request_data["source_name"] = path_obj.name zip_path = self._zip_directory(path) @@ -1847,12 +1855,14 @@ class SyncHTTPClient: args: Optional[Dict[str, Any]] = None, telemetry: Any = False, processing_mode: Optional[str] = None, + add_type: Optional[str] = None, tags: Optional[List[str]] = None, tag_mode: str = "replace", ) -> Dict[str, Any]: return run_async( self._async_client.add_resource( path=path, + add_type=add_type, to=to, parent=parent, reason=reason, diff --git a/sdk/python/tests/test_async_client_behaviors.py b/sdk/python/tests/test_async_client_behaviors.py index fa76c484d..bc374e48a 100644 --- a/sdk/python/tests/test_async_client_behaviors.py +++ b/sdk/python/tests/test_async_client_behaviors.py @@ -480,6 +480,93 @@ async def test_add_resource_forwards_processing_mode(): assert payload["processing_mode"] == "vectors_only" +@pytest.mark.asyncio +async def test_add_resource_forwards_declared_add_type_with_exact_target(): + client = AsyncHTTPClient(url="http://127.0.0.1:1933") + fake_http = SimpleNamespace(post=AsyncMock(return_value=object())) + client._http = fake_http + client._handle_response_data = lambda _response: { + "result": {"root_uri": "viking://resources/feishu"} + } + + await client.add_resource( + "space:home", + add_type=" feishu ", + to="viking://resources/feishu", + ) + + payload = fake_http.post.await_args.kwargs["json"] + assert payload["path"] == "space:home" + assert payload["add_type"] == "feishu" + assert payload["to"] == "viking://resources/feishu" + + +@pytest.mark.asyncio +async def test_add_resource_declared_add_type_requires_exact_target(): + client = AsyncHTTPClient(url="http://127.0.0.1:1933") + + with pytest.raises(ValueError, match="exact 'to'"): + await client.add_resource("space:home", add_type="feishu") + + +@pytest.mark.asyncio +async def test_add_resource_declared_add_type_rejects_parent(): + client = AsyncHTTPClient(url="http://127.0.0.1:1933") + + with pytest.raises(ValueError, match="'parent'"): + await client.add_resource( + "space:home", + add_type="feishu", + to="viking://resources/feishu", + parent="viking://resources/imports", + ) + + +@pytest.mark.asyncio +async def test_add_resource_declared_add_type_skips_local_file_upload(tmp_path): + source = tmp_path / "source" + source.write_text("connector source") + + client = AsyncHTTPClient(url="http://127.0.0.1:1933") + fake_http = SimpleNamespace(post=AsyncMock(return_value=object())) + client._http = fake_http + client._upload_temp_file = AsyncMock(return_value="unexpected-upload") + client._handle_response_data = lambda _response: { + "result": {"root_uri": "viking://resources/feishu"} + } + + await client.add_resource( + str(source), + add_type="feishu", + to="viking://resources/feishu", + ) + + client._upload_temp_file.assert_not_awaited() + payload = fake_http.post.await_args.kwargs["json"] + assert payload["path"] == str(source) + assert "temp_file_id" not in payload + + +def test_sync_add_resource_accepts_and_forwards_declared_add_type(): + client = SyncHTTPClient(url="http://127.0.0.1:1933") + + with patch.object( + client._async_client, + "add_resource", + new_callable=AsyncMock, + return_value={"root_uri": "viking://resources/feishu"}, + ) as mock_add_resource: + result = client.add_resource( + "space:home", + add_type="feishu", + to="viking://resources/feishu", + ) + + assert result["root_uri"] == "viking://resources/feishu" + assert mock_add_resource.await_args.kwargs["add_type"] == "feishu" + assert mock_add_resource.await_args.kwargs["to"] == "viking://resources/feishu" + + @pytest.mark.asyncio async def test_add_resource_omits_default_processing_mode_for_legacy_servers(): client = AsyncHTTPClient(url="http://127.0.0.1:1933") diff --git a/tests/client/test_add_resource_signature_compat.py b/tests/client/test_add_resource_signature_compat.py index 76185fe2c..a339abc8d 100644 --- a/tests/client/test_add_resource_signature_compat.py +++ b/tests/client/test_add_resource_signature_compat.py @@ -1,8 +1,21 @@ import inspect +from openviking_sdk.client import AsyncHTTPClient, SyncHTTPClient + from openviking import AsyncOpenViking, SyncOpenViking from openviking.client.local import LocalClient -from openviking_sdk.client import AsyncHTTPClient, SyncHTTPClient + + +def test_python_add_resource_clients_accept_add_type_keyword(): + for client_type in ( + AsyncOpenViking, + SyncOpenViking, + LocalClient, + AsyncHTTPClient, + SyncHTTPClient, + ): + parameters = inspect.signature(client_type.add_resource).parameters + assert "add_type" in parameters, client_type.__name__ def test_async_openviking_add_resource_preserves_positional_watch_args(): diff --git a/tests/client/test_http_client_compact.py b/tests/client/test_http_client_compact.py index cdedaa4da..547b4e289 100644 --- a/tests/client/test_http_client_compact.py +++ b/tests/client/test_http_client_compact.py @@ -103,7 +103,15 @@ async def test_add_resource_omits_empty_args_and_null_fields(): assert payload["path"] == "https://example.com/doc" # `args` is the field that breaks `add-resource` against pre-#2549 instances. assert "args" not in payload - for dropped in ("to", "parent", "timeout", "ignore_dirs", "include", "exclude"): + for dropped in ( + "add_type", + "to", + "parent", + "timeout", + "ignore_dirs", + "include", + "exclude", + ): assert dropped not in payload diff --git a/tests/client/test_resource_management.py b/tests/client/test_resource_management.py index 1015c641b..09cdbd261 100644 --- a/tests/client/test_resource_management.py +++ b/tests/client/test_resource_management.py @@ -76,6 +76,48 @@ class TestAddResource: assert str(seen["telemetry_id"]).startswith("tm_") assert seen["kwargs"]["wait"] is True + async def test_local_client_forwards_declared_add_type(self): + seen: dict[str, object] = {} + + async def _fake_add_resource(**kwargs): + seen.update(kwargs) + return {"root_uri": "viking://resources/feishu"} + + client = LocalClient.__new__(LocalClient) + client._ctx = RequestContext(user=UserIdentifier.the_default_user(), role=Role.USER) + client._service = SimpleNamespace( + resources=SimpleNamespace(add_resource=_fake_add_resource) + ) + + result = await LocalClient.add_resource( + client, + path="space:home", + add_type=" feishu ", + to="viking://resources/feishu", + ) + + assert result["root_uri"] == "viking://resources/feishu" + assert seen["path"] == "space:home" + assert seen["add_type"] == "feishu" + assert seen["to"] == "viking://resources/feishu" + + async def test_async_openviking_forwards_declared_add_type(self): + backend = SimpleNamespace(add_resource=AsyncMock(return_value={"root_uri": "ok"})) + client = AsyncOpenViking.__new__(AsyncOpenViking) + client._initialized = True + client._client = backend + + result = await AsyncOpenViking.add_resource( + client, + path="space:home", + add_type="feishu", + to="viking://resources/feishu", + ) + + assert result == {"root_uri": "ok"} + assert backend.add_resource.await_args.kwargs["add_type"] == "feishu" + assert backend.add_resource.await_args.kwargs["to"] == "viking://resources/feishu" + async def test_add_resource_without_wait( self, client: AsyncOpenViking, sample_markdown_file: Path ): diff --git a/tests/connector/test_client.py b/tests/connector/test_client.py index c2bb6519a..501cbb688 100644 --- a/tests/connector/test_client.py +++ b/tests/connector/test_client.py @@ -103,6 +103,27 @@ async def test_submit_doc_add_carries_non_tos_source_in_param_config(): assert "tos_path" not in payload +@pytest.mark.asyncio +async def test_submit_doc_add_forwards_resource_tags(): + _FakeAsyncClient.response_payload = {"code": 0, "data": {"task_key": "connector-1"}} + client = ConnectorClient("https://connector/doc/add", "https://tracker/task/info", "acct") + + await client.submit_doc_add( + add_type="tos", + api_key="secret", + tos_path="bucket/prefix", + to="viking://resources/imports", + extra_params={ + "tags": ["team=search", "env=test"], + "tag_mode": "append", + }, + ) + + payload = _FakeAsyncClient.calls[0]["json"] + assert payload["tags"] == ["team=search", "env=test"] + assert payload["tag_mode"] == "append" + + @pytest.mark.asyncio async def test_submit_doc_add_keeps_credentials_out_of_param_config(): _FakeAsyncClient.response_payload = {"code": 0, "data": {"task_key": "connector-1"}} diff --git a/tests/parse/test_feishu_parser_api.py b/tests/parse/test_feishu_parser_api.py index 4dbbded75..22a732b34 100644 --- a/tests/parse/test_feishu_parser_api.py +++ b/tests/parse/test_feishu_parser_api.py @@ -450,7 +450,11 @@ async def test_uat_producer_payload_reaches_worker_without_persisting_token(monk resource_processor=resource_processor, skill_processor=SimpleNamespace(), ) - monkeypatch.setattr(service, "_should_use_connector", Mock(return_value=False)) + monkeypatch.setattr( + service, + "_connector_delegate", + SimpleNamespace(should_delegate=Mock(return_value=False)), + ) monkeypatch.setattr( "openviking.service.resource_service.is_git_repo_url", Mock(return_value=False), @@ -527,7 +531,11 @@ async def test_local_prepared_job_uses_add_resource_queue(monkeypatch): skill_processor=SimpleNamespace(), ) service._enqueue_add_resource_job = AsyncMock(return_value=SimpleNamespace(task_id="task-1")) - monkeypatch.setattr(service, "_should_use_connector", Mock(return_value=False)) + monkeypatch.setattr( + service, + "_connector_delegate", + SimpleNamespace(should_delegate=Mock(return_value=False)), + ) monkeypatch.setattr( "openviking.service.resource_service.is_git_repo_url", Mock(return_value=False), @@ -615,7 +623,11 @@ async def test_uat_producer_cancellation_respects_queue_ownership( resource_processor=resource_processor, skill_processor=SimpleNamespace(), ) - monkeypatch.setattr(service, "_should_use_connector", Mock(return_value=False)) + monkeypatch.setattr( + service, + "_connector_delegate", + SimpleNamespace(should_delegate=Mock(return_value=False)), + ) monkeypatch.setattr( "openviking.service.resource_service.is_git_repo_url", Mock(return_value=False), diff --git a/tests/server/test_api_resources.py b/tests/server/test_api_resources.py index 5ac3312fa..9ddecf0cc 100644 --- a/tests/server/test_api_resources.py +++ b/tests/server/test_api_resources.py @@ -33,6 +33,62 @@ def test_add_resource_request_defaults_processing_mode(): assert request.processing_mode == "semantic_and_vectors" +def test_add_resource_request_accepts_declared_add_type(): + request = AddResourceRequest( + path="https://example.com/space", + add_type=" feishu ", + to="viking://resources/feishu", + ) + + assert request.add_type == "feishu" + assert request.to == "viking://resources/feishu" + + +def test_add_resource_request_rejects_add_type_with_temp_file_id(): + import pytest + + with pytest.raises(ValueError, match="temp_file_id"): + AddResourceRequest(temp_file_id="upload_abc", add_type="feishu") + + +def test_add_resource_request_requires_exact_to_for_declared_add_type(): + import pytest + + with pytest.raises(ValueError, match="exact 'to'"): + AddResourceRequest(path="space:home", add_type="feishu") + + +def test_add_resource_request_rejects_add_type_with_parent(): + import pytest + + with pytest.raises(ValueError, match="'parent'"): + AddResourceRequest( + path="space:home", + add_type="feishu", + to="viking://resources/feishu", + parent="viking://resources/imports", + ) + + +def test_require_remote_resource_source_allows_declared_add_type(): + from openviking.server.local_input_guard import require_remote_resource_source + + assert ( + require_remote_resource_source("space:home", declared_connector_add_type="feishu") + == "space:home" + ) + + +def test_require_remote_resource_source_still_rejects_without_declared_type(): + import pytest + + from openviking.server.local_input_guard import require_remote_resource_source + from openviking_cli.exceptions import PermissionDeniedError + + with pytest.raises(PermissionDeniedError): + require_remote_resource_source("/etc/passwd") + + async def _wait_task_terminal(client: httpx.AsyncClient, task_id: str, timeout: float = 10.0): deadline = asyncio.get_running_loop().time() + timeout last_result = None diff --git a/tests/server/test_mcp_endpoint.py b/tests/server/test_mcp_endpoint.py index abfb45730..27ab9d81f 100644 --- a/tests/server/test_mcp_endpoint.py +++ b/tests/server/test_mcp_endpoint.py @@ -643,6 +643,53 @@ async def test_add_resource_remote_parent_is_forwarded(service, monkeypatch): assert captured["to"] is None +async def test_add_resource_declared_add_type_is_forwarded(service, monkeypatch): + captured = {} + + async def fake_add_resource(*, path, ctx, **kwargs): + captured["path"] = path + captured.update(kwargs) + return {"task_id": "ov-task-123"} + + monkeypatch.setattr(service.resources, "add_resource", fake_add_resource) + + # A non-URL source: only reaches the service because add_type is declared; + # without it this path shape would be treated as a local file. + result = await add_resource( + path="space:home", + add_type="feishu", + to="viking://resources/feishu", + ) + + assert "task_id: ov-task-123" in result + assert captured["path"] == "space:home" + assert captured["add_type"] == "feishu" + assert captured["to"] == "viking://resources/feishu" + + +async def test_add_resource_declared_add_type_rejects_temp_file_id(service): + result = await add_resource(temp_file_id="upload_abc.md", add_type="feishu") + + assert result == "Error: add_type cannot be combined with temp_file_id." + + +async def test_add_resource_declared_add_type_requires_exact_to(service): + result = await add_resource(path="space:home", add_type="feishu") + + assert result == "Error: add_type requires an exact 'to' target." + + +async def test_add_resource_declared_add_type_rejects_parent(service): + result = await add_resource( + path="space:home", + add_type="feishu", + to="viking://resources/feishu", + parent="viking://resources/imports", + ) + + assert result == "Error: add_type cannot be combined with parent." + + async def test_add_resource_remote_tags_are_forwarded(service, monkeypatch): captured = {} diff --git a/tests/service/test_resource_service_connector.py b/tests/service/test_resource_service_connector.py index 3fca0863c..d9cec7699 100644 --- a/tests/service/test_resource_service_connector.py +++ b/tests/service/test_resource_service_connector.py @@ -9,6 +9,7 @@ from unittest.mock import AsyncMock, Mock import httpx import pytest +from openviking.connector import delegate as connector_delegate_module from openviking.server.identity import RequestContext, Role from openviking.service import resource_service as resource_service_module from openviking.service.resource_service import ResourceService @@ -100,7 +101,7 @@ def _install_connector_dependencies(monkeypatch, tracker, connector_client): lambda: tracker, ) monkeypatch.setattr( - resource_service_module, + connector_delegate_module, "ConnectorClient", lambda **_kwargs: connector_client, ) @@ -109,7 +110,7 @@ def _install_connector_dependencies(monkeypatch, tracker, connector_client): coro.close() return _BackgroundTask() - monkeypatch.setattr(resource_service_module.asyncio, "create_task", discard_monitor) + monkeypatch.setattr(connector_delegate_module.asyncio, "create_task", discard_monitor) @pytest.mark.asyncio @@ -260,7 +261,7 @@ def test_git_commit_tree_url_falls_back_to_native(connector_config, service): connector_config.allowed_add_types = ["git"] assert ( - service._should_use_connector( + service._connector.should_delegate( "https://github.com/acme/repo/tree/deadbee", to="viking://resources/imports", ) @@ -272,7 +273,7 @@ def test_git_commit_tree_url_with_credentials_fails_closed(connector_config, ser connector_config.allowed_add_types = ["git"] with pytest.raises(InvalidArgumentError, match="cannot fall back") as exc_info: - service._should_use_connector( + service._connector.should_delegate( "https://github.com/acme/repo/tree/deadbee", to="viking://resources/imports", connector_args={"token": "ghp-secret"}, @@ -288,7 +289,7 @@ def test_git_full_commit_arg_routes_to_connector(connector_config, service): connector_config.allowed_add_types = ["git"] assert ( - service._should_use_connector( + service._connector.should_delegate( "https://git.example/org/repo.git", to="viking://resources/imports", connector_args={"commit": _FULL_COMMIT_SHA}, @@ -301,7 +302,7 @@ def test_git_full_commit_with_credentials_routes_to_connector(connector_config, connector_config.allowed_add_types = ["git"] assert ( - service._should_use_connector( + service._connector.should_delegate( "https://git.example/org/repo.git", to="viking://resources/imports", connector_args={"commit": _FULL_COMMIT_SHA, "token": "ghp-secret"}, @@ -319,7 +320,7 @@ def test_git_commit_with_branch_falls_back_to_native_for_wait( connector_config.allowed_add_types = ["git"] assert ( - service._should_use_connector( + service._connector.should_delegate( "https://git.example/org/repo.git", to="viking://resources/imports", wait=True, @@ -483,23 +484,23 @@ def test_connector_only_route_rejects_disabled_or_unsupported_requests( connector_config, service, ): - assert service._should_use_connector("https://example.com/doc") is False + assert service._connector.should_delegate("https://example.com/doc") is False with pytest.raises(InvalidArgumentError, match="args keys"): - service._should_use_connector("tos://bucket/prefix", connector_args={"parser": "pdf"}) + service._connector.should_delegate("tos://bucket/prefix", connector_args={"parser": "pdf"}) connector_config.enable = False with pytest.raises(InvalidArgumentError, match="Connector integration"): - service._should_use_connector("tos://bucket/prefix") + service._connector.should_delegate("tos://bucket/prefix") def test_git_route_degrades_when_disabled_or_type_not_allowed(connector_config, service): # "git" not in allowed_add_types: standard pipeline handles the repo. - assert service._should_use_connector("https://git.example/org/repo.git") is False + assert service._connector.should_delegate("https://git.example/org/repo.git") is False connector_config.allowed_add_types = ["tos", "git"] connector_config.enable = False - assert service._should_use_connector("https://git.example/org/repo.git") is False + assert service._connector.should_delegate("https://git.example/org/repo.git") is False @pytest.mark.parametrize( @@ -563,7 +564,7 @@ def test_git_source_falls_back_for_parent_target( connector_config.allowed_add_types = ["tos", "git"] assert ( - service._should_use_connector( + service._connector.should_delegate( "https://git.example/org/repo.git", parent="viking://resources/manuals", ) @@ -575,7 +576,7 @@ def test_git_route_accepts_explicit_create_parent_false(connector_config, servic connector_config.allowed_add_types = ["tos", "git"] assert ( - service._should_use_connector( + service._connector.should_delegate( "https://git.example/org/repo.git", to="viking://resources/repo", kwargs={"create_parent": False}, @@ -588,7 +589,7 @@ def test_git_source_falls_back_for_unsupported_args(connector_config, service): connector_config.allowed_add_types = ["tos", "git"] assert ( - service._should_use_connector( + service._connector.should_delegate( "https://git.example/org/repo.git", connector_args={"depth": "1"}, ) @@ -600,7 +601,7 @@ def test_git_route_accepts_credential_args(connector_config, service): connector_config.allowed_add_types = ["tos", "git"] assert ( - service._should_use_connector( + service._connector.should_delegate( "https://git.example/org/repo.git", to="viking://resources/repo", connector_args={"token": "ghp-secret", "username": "oauth2"}, @@ -627,7 +628,7 @@ def test_git_credentials_fail_closed_when_connector_request_would_fallback( connector_config.allowed_add_types = ["tos", "git"] with pytest.raises(InvalidArgumentError, match="cannot fall back") as exc_info: - service._should_use_connector( + service._connector.should_delegate( "https://git.example/org/private.git", connector_args={"token": "ghp-secret", "username": "oauth2"}, **routing_kwargs, @@ -638,7 +639,7 @@ def test_git_credentials_fail_closed_when_connector_request_would_fallback( def test_git_credentials_require_enabled_allowed_connector(connector_config, service): with pytest.raises(InvalidArgumentError, match="require Connector import") as exc_info: - service._should_use_connector( + service._connector.should_delegate( "https://git.example/org/private.git", to="viking://resources/private", connector_args={"token": "ghp-secret"}, @@ -649,7 +650,7 @@ def test_git_credentials_require_enabled_allowed_connector(connector_config, ser def test_tos_route_rejects_credential_args(connector_config, service): with pytest.raises(InvalidArgumentError, match="args keys"): - service._should_use_connector( + service._connector.should_delegate( "tos://bucket/prefix", connector_args={"token": "ghp-secret"}, ) @@ -693,7 +694,7 @@ async def test_git_credentials_with_include_never_enqueue_native_job( ): connector_config.allowed_add_types = ["tos", "git"] monkeypatch.setattr(resource_service_module, "is_git_repo_url", lambda _path: True) - service._add_resource_via_connector = AsyncMock() + service._connector.submit = AsyncMock() service.enqueue_git_add_resource = AsyncMock() with pytest.raises(InvalidArgumentError, match="cannot fall back") as exc_info: @@ -707,7 +708,7 @@ async def test_git_credentials_with_include_never_enqueue_native_job( ) assert "ghp-secret" not in str(exc_info.value) - service._add_resource_via_connector.assert_not_awaited() + service._connector.submit.assert_not_awaited() service.enqueue_git_add_resource.assert_not_awaited() @@ -744,7 +745,7 @@ def test_connector_route_accepts_public_exact_to(connector_config, ctx, service, connector_config.allowed_add_types = ["tos", "git"] assert ( - service._should_use_connector( + service._connector.should_delegate( "https://git.example/org/repo.git", ctx=ctx, to=to, @@ -762,7 +763,7 @@ async def test_add_resource_falls_back_for_shared_source_with_parent( ): connector_config.allowed_add_types = ["tos", "git"] monkeypatch.setattr(resource_service_module, "is_git_repo_url", lambda _path: True) - service._add_resource_via_connector = AsyncMock() + service._connector.submit = AsyncMock() service.enqueue_git_add_resource = AsyncMock(return_value={"root_uri": "standard-pipeline"}) result = await service.add_resource( @@ -772,7 +773,7 @@ async def test_add_resource_falls_back_for_shared_source_with_parent( ) assert result == {"root_uri": "standard-pipeline"} - service._add_resource_via_connector.assert_not_awaited() + service._connector.submit.assert_not_awaited() service.enqueue_git_add_resource.assert_awaited_once() @@ -845,7 +846,7 @@ async def test_shared_source_create_parent_false_routes_to_connector( ): connector_config.allowed_add_types = ["tos", "git"] monkeypatch.setattr(resource_service_module, "is_git_repo_url", lambda _path: True) - service._add_resource_via_connector = AsyncMock(return_value={"status": "accepted"}) + service._connector.submit = AsyncMock(return_value={"status": "accepted"}) service.enqueue_git_add_resource = AsyncMock() result = await service.add_resource( @@ -857,7 +858,7 @@ async def test_shared_source_create_parent_false_routes_to_connector( assert result == {"status": "accepted"} service.enqueue_git_add_resource.assert_not_awaited() - service._add_resource_via_connector.assert_awaited_once() + service._connector.submit.assert_awaited_once() @pytest.mark.asyncio @@ -874,7 +875,7 @@ async def test_connector_import_without_target_is_rejected( _install_connector_dependencies(monkeypatch, tracker, connector_client) with pytest.raises(InvalidArgumentError, match="exact 'to' target"): - await service._add_resource_via_connector( + await service._connector.submit( path="tos://bucket/prefix", ctx=ctx, to=None, @@ -916,7 +917,7 @@ def test_git_source_falls_back_for_wait(connector_config, service): connector_config.allowed_add_types = ["tos", "git"] assert ( - service._should_use_connector( + service._connector.should_delegate( "https://git.example/org/repo.git", to="viking://resources/repo", wait=True, @@ -937,14 +938,31 @@ async def test_tos_connector_rejects_wait(connector_config, ctx, service): @pytest.mark.asyncio -async def test_tos_connector_rejects_tags_instead_of_dropping_them(connector_config, ctx, service): - with pytest.raises(InvalidArgumentError, match="tags"): - await service.add_resource( - path="tos://bucket/prefix", - ctx=ctx, - to="viking://resources/imports", - tags=["team=search"], - ) +async def test_tos_connector_forwards_tags_and_mode( + monkeypatch, + connector_config, + ctx, + service, +): + tracker = _task_tracker() + connector_client = SimpleNamespace( + submit_doc_add=AsyncMock(return_value={"task_key": "connector-1"}) + ) + _install_connector_dependencies(monkeypatch, tracker, connector_client) + + await service.add_resource( + path="tos://bucket/prefix", + ctx=ctx, + to="viking://resources/imports", + tags=[" Team=Search ", "env=test", "team=search"], + tag_mode="append", + ) + + submitted = connector_client.submit_doc_add.await_args.kwargs + assert submitted["extra_params"] == { + "tags": ["team=search", "env=test"], + "tag_mode": "append", + } @pytest.mark.asyncio @@ -962,7 +980,7 @@ async def test_monitor_links_reason_memory_on_success( client = SimpleNamespace(get_task_info=AsyncMock(return_value={"Status": "succeeded"})) service._link_resource_reason_memory = AsyncMock() - outcome = await service._monitor_connector_task( + outcome = await service._connector._monitor( client=client, connector_task_key="connector-1", ov_task_id="task-1", @@ -989,7 +1007,7 @@ async def test_git_reason_routes_to_connector( ): connector_config.allowed_add_types = ["tos", "git"] monkeypatch.setattr(resource_service_module, "is_git_repo_url", lambda _path: True) - service._add_resource_via_connector = AsyncMock(return_value={"status": "accepted"}) + service._connector.submit = AsyncMock(return_value={"status": "accepted"}) service.enqueue_git_add_resource = AsyncMock() result = await service.add_resource( @@ -1001,7 +1019,7 @@ async def test_git_reason_routes_to_connector( assert result == {"status": "accepted"} service.enqueue_git_add_resource.assert_not_awaited() - connector_kwargs = service._add_resource_via_connector.await_args.kwargs + connector_kwargs = service._connector.submit.await_args.kwargs assert connector_kwargs["reason"] == "track quarterly reports" @@ -1034,10 +1052,10 @@ async def test_monitor_connector_task_maps_terminal_status( async def no_sleep(_seconds): pass - monkeypatch.setattr(resource_service_module.asyncio, "sleep", no_sleep) + monkeypatch.setattr(connector_delegate_module.asyncio, "sleep", no_sleep) client = SimpleNamespace(get_task_info=AsyncMock(return_value=task_info)) - outcome = await ResourceService()._monitor_connector_task( + outcome = await ResourceService()._connector._monitor( client=client, connector_task_key="connector-1", ov_task_id="task-1", @@ -1072,7 +1090,7 @@ async def test_monitor_connector_task_retries_transient_polling_error( async def no_sleep(_seconds): pass - monkeypatch.setattr(resource_service_module.asyncio, "sleep", no_sleep) + monkeypatch.setattr(connector_delegate_module.asyncio, "sleep", no_sleep) client = SimpleNamespace( get_task_info=AsyncMock( side_effect=[ @@ -1082,7 +1100,7 @@ async def test_monitor_connector_task_retries_transient_polling_error( ) ) - await ResourceService()._monitor_connector_task( + await ResourceService()._connector._monitor( client=client, connector_task_key="connector-1", ov_task_id="task-1", @@ -1111,11 +1129,11 @@ async def test_monitor_connector_task_marks_cancelled_monitor_as_failed( async def cancelled_sleep(_seconds): raise asyncio.CancelledError - monkeypatch.setattr(resource_service_module.asyncio, "sleep", cancelled_sleep) + monkeypatch.setattr(connector_delegate_module.asyncio, "sleep", cancelled_sleep) client = SimpleNamespace(get_task_info=AsyncMock()) with pytest.raises(asyncio.CancelledError): - await ResourceService()._monitor_connector_task( + await ResourceService()._connector._monitor( client=client, connector_task_key="connector-1", ov_task_id="task-1", @@ -1131,3 +1149,185 @@ async def test_monitor_connector_task_marks_cancelled_monitor_as_failed( user_id="alice", ) tracker.complete.assert_not_awaited() + + +# --- Declared add_type routing ------------------------------------------------- + + +@pytest.mark.asyncio +async def test_declared_add_type_routes_generic_type_to_connector( + monkeypatch, + connector_config, + ctx, + service, +): + connector_config.allowed_add_types = ["feishu"] + tracker = _task_tracker() + connector_client = SimpleNamespace( + submit_doc_add=AsyncMock(return_value={"task_key": "connector-1"}) + ) + _install_connector_dependencies(monkeypatch, tracker, connector_client) + + result = await service.add_resource( + path="https://example.feishu.cn/wiki/space-home", + ctx=ctx, + add_type="feishu", + to="viking://resources/kb/feishu", + args={"space_id": "spc1", "auth_config": {"app_secret": "s3cret"}}, + ) + + assert result["status"] == "accepted" + connector_client.submit_doc_add.assert_awaited_once_with( + add_type="feishu", + api_key="secret", + tos_path=None, + to="viking://resources/kb/feishu", + include_child=True, + param_config={ + "space_id": "spc1", + "path": "https://example.feishu.cn/wiki/space-home", + }, + auth_config={"app_secret": "s3cret"}, + extra_params=None, + ) + + +@pytest.mark.asyncio +async def test_declared_registry_type_routes_like_probed( + monkeypatch, + connector_config, + ctx, + service, +): + tracker = _task_tracker() + connector_client = SimpleNamespace( + submit_doc_add=AsyncMock(return_value={"task_key": "connector-1"}) + ) + _install_connector_dependencies(monkeypatch, tracker, connector_client) + + await service.add_resource( + path="tos://bucket/a/b", + ctx=ctx, + add_type="tos", + to="viking://resources/x", + ) + + connector_client.submit_doc_add.assert_awaited_once_with( + add_type="tos", + api_key="secret", + tos_path="bucket/a/b", + to="viking://resources/x", + include_child=True, + param_config=None, + auth_config=None, + extra_params=None, + ) + + +def test_declared_add_type_requires_enabled_and_allowed(connector_config, ctx, service): + # Not in allowed_add_types (fixture allows only "tos"). + with pytest.raises(InvalidArgumentError, match="disabled or does not allow"): + service._connector.should_delegate( + "https://example.feishu.cn/wiki/x", + ctx=ctx, + declared_add_type="feishu", + to="viking://resources/x", + ) + # Connector disabled entirely. + connector_config.enable = False + connector_config.allowed_add_types = ["feishu"] + with pytest.raises(InvalidArgumentError, match="disabled or does not allow"): + service._connector.should_delegate( + "https://example.feishu.cn/wiki/x", + ctx=ctx, + declared_add_type="feishu", + to="viking://resources/x", + ) + + +def test_declared_add_type_rejects_probe_mismatch(connector_config, ctx, service): + connector_config.allowed_add_types = ["tos", "git", "feishu"] + # Generic declared type claiming a path that probes as a registry type. + with pytest.raises(InvalidArgumentError, match="does not match the source path"): + service._connector.should_delegate( + "tos://bucket/x", + ctx=ctx, + declared_add_type="feishu", + to="viking://resources/x", + ) + # Registry declared type on a path that probes differently (or not at all). + with pytest.raises(InvalidArgumentError, match="does not match the source path"): + service._connector.should_delegate( + "https://not-a-repo.example/page", + ctx=ctx, + declared_add_type="git", + to="viking://resources/x", + ) + with pytest.raises(InvalidArgumentError, match="does not match the source path"): + service._connector.should_delegate( + _GIT_REPO_PREFIX + "org/repo", + ctx=ctx, + declared_add_type="tos", + to="viking://resources/x", + ) + + +def test_declared_add_type_never_degrades(connector_config, ctx, service): + connector_config.allowed_add_types = ["tos", "git"] + # Probed git with wait=true degrades; declared git must raise instead. + with pytest.raises(InvalidArgumentError, match="wait=true"): + service._connector.should_delegate( + _GIT_REPO_PREFIX + "org/repo", + ctx=ctx, + declared_add_type="git", + to="viking://resources/x", + wait=True, + ) + # Same for an unsupported top-level param on a generic declared type. + connector_config.allowed_add_types = ["feishu"] + with pytest.raises(InvalidArgumentError, match="instruction"): + service._connector.should_delegate( + "https://example.feishu.cn/wiki/x", + ctx=ctx, + declared_add_type="feishu", + to="viking://resources/x", + instruction="summarize", + ) + + +def test_declared_git_rejects_abbreviated_commit(connector_config, ctx, service): + connector_config.allowed_add_types = ["git"] + with pytest.raises(InvalidArgumentError, match="abbreviated commit SHA"): + service._connector.should_delegate( + _GIT_REPO_PREFIX + "org/repo", + ctx=ctx, + declared_add_type="git", + to="viking://resources/x", + connector_args={"commit": "abc1234"}, + ) + + +def test_declared_add_type_rejects_non_mapping_auth_config(connector_config, ctx, service): + connector_config.allowed_add_types = ["feishu"] + with pytest.raises(InvalidArgumentError, match="auth_config"): + service._connector.should_delegate( + "https://example.feishu.cn/wiki/x", + ctx=ctx, + declared_add_type="feishu", + to="viking://resources/x", + connector_args={"auth_config": "not-a-dict"}, + ) + + +def test_declared_git_keeps_args_whitelist(connector_config, ctx, service): + # Registry types keep their per-type args whitelist even when declared; + # unknown keys raise instead of flowing into param_config. + connector_config.allowed_add_types = ["git"] + with pytest.raises(InvalidArgumentError, match="args keys"): + service._connector.should_delegate( + _GIT_REPO_PREFIX + "org/repo", + ctx=ctx, + declared_add_type="git", + to="viking://resources/x", + connector_args={"depth": 1}, + ) diff --git a/tests/service/test_resource_service_understanding_routing.py b/tests/service/test_resource_service_understanding_routing.py index 7bb048440..25ccbfe96 100644 --- a/tests/service/test_resource_service_understanding_routing.py +++ b/tests/service/test_resource_service_understanding_routing.py @@ -70,7 +70,7 @@ async def test_extensionless_remote_url_queues_frozen_understanding_route( resource_processor=processor, skill_processor=object(), ) - service._should_use_connector = lambda *_args, **_kwargs: False + service._connector_delegate = SimpleNamespace(should_delegate=lambda *_args, **_kwargs: False) tracker = SimpleNamespace( create=AsyncMock(return_value=SimpleNamespace(task_id="task-1")), update_stage=AsyncMock(),