refactor: extract Connector delegation and expose declarative add_type (#3591)

* refactor: delegate add_resource imports to external Connector

* refactor: delegate add_resource imports to external Connector

* refactor: extract Connector delegation and expose declarative add_type

* fix: merge main to refactor/connector_delegator

* fix: merge main to refactor/connector_delegator
This commit is contained in:
zihengli
2026-07-29 18:09:25 +08:00
committed by GitHub
parent 291e9c580f
commit e9c4cc97c3
27 changed files with 1412 additions and 562 deletions
+5 -1
View File
@@ -692,6 +692,7 @@ impl HttpClient {
pub async fn add_resource(
&self,
path: &str,
add_type: Option<String>,
to: Option<String>,
parent: Option<String>,
parent_auto_create: Option<String>,
@@ -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,
+2
View File
@@ -6,6 +6,7 @@ use serde_json::{Map, Value};
pub async fn add_resource(
client: &HttpClient,
path: &str,
add_type: Option<String>,
to: Option<String>,
parent: Option<String>,
parent_auto_create: Option<String>,
@@ -31,6 +32,7 @@ pub async fn add_resource(
let result = client
.add_resource(
path,
add_type,
to,
parent,
parent_auto_create,
+5 -1
View File
@@ -15,6 +15,7 @@ use serde_json::{Map, Value};
pub async fn handle_add_resource(
mut path: String,
add_type: Option<String>,
to: Option<String>,
parent: Option<String>,
parent_auto_create: Option<String>,
@@ -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,
+4
View File
@@ -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 {
+69
View File
@@ -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<String>,
/// Exact target URI (must not exist yet) (cannot be used with --parent)
#[arg(long, value_name = "uri", help_heading = "Common options")]
to: Option<String>,
@@ -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([
+11
View File
@@ -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,
+13 -1
View File
@@ -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,
+3 -2
View File
@@ -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).
"""
+612
View File
@@ -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://<bucket>/<path>'."
)
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)}
+21 -4
View File
@@ -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.
+11 -3
View File
@@ -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."
+30 -6
View File
@@ -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
+26 -3
View File
@@ -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,
+35 -488
View File
@@ -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://<bucket>/<path>'."
)
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,
+11
View File
@@ -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,
+6 -1
View File
@@ -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
+11 -1
View File
@@ -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,
@@ -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")
@@ -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():
+9 -1
View File
@@ -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
+42
View File
@@ -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
):
+21
View File
@@ -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"}}
+15 -3
View File
@@ -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),
+56
View File
@@ -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
+47
View File
@@ -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 = {}
+245 -45
View File
@@ -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},
)
@@ -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(),