# Copyright (c) 2026 Beijing Volcano Engine Technology Co., Ltd. # SPDX-License-Identifier: AGPL-3.0 """ Async OpenViking client implementation (embedded mode only). For HTTP mode, use AsyncHTTPClient or SyncHTTPClient. """ from __future__ import annotations import threading from typing import TYPE_CHECKING, Any, Dict, List, Optional, Union from openviking.client import LocalClient, Session from openviking.service.debug_service import SystemStatus from openviking.telemetry import TelemetryRequest from openviking.utils.search_filters import SearchContextTypeInput from openviking_cli.client.base import BaseClient from openviking_cli.session.user_id import UserIdentifier from openviking_cli.utils import get_logger if TYPE_CHECKING: from openviking.snapshot_namespace import AsyncSnapshotNamespace logger = get_logger(__name__) if TYPE_CHECKING: from openviking.snapshot_namespace import AsyncSnapshotNamespace class AsyncOpenViking: """ OpenViking main client class (Asynchronous, embedded mode only). Uses local storage and auto-starts services (singleton). For HTTP mode, use AsyncHTTPClient or SyncHTTPClient instead. Examples: client = AsyncOpenViking(path="./data") await client.initialize() """ _instance: Optional["AsyncOpenViking"] = None _lock = threading.Lock() def __new__(cls, *args, **kwargs): if cls._instance is None: with cls._lock: if cls._instance is None: cls._instance = object.__new__(cls) return cls._instance def __init__( self, path: Optional[str] = None, actor_peer_id: Optional[str] = None, agent_id: Optional[str] = None, ): """ Initialize OpenViking client (embedded mode). Args: path: Local storage path (overrides ov.conf storage path). actor_peer_id: Optional view filter for the current user's peer collection. agent_id: Legacy alias for actor_peer_id. """ # Singleton guard for repeated initialization if hasattr(self, "_singleton_initialized") and self._singleton_initialized: return self.user = UserIdentifier.the_default_user() self._initialized = False self._snapshot: Optional["AsyncSnapshotNamespace"] = None # Mark initialized only after LocalClient is successfully constructed. self._singleton_initialized = False self._client: BaseClient = LocalClient( path=path, actor_peer_id=actor_peer_id, agent_id=agent_id, ) self._singleton_initialized = True # ============= Lifecycle methods ============= async def initialize(self) -> None: """Initialize OpenViking storage and indexes.""" await self._client.initialize() self._initialized = True async def _ensure_initialized(self): """Ensure storage collections are initialized.""" if not self._initialized: await self.initialize() async def close(self) -> None: """Close OpenViking and release resources.""" client = getattr(self, "_client", None) if client is not None: await client.close() self._initialized = False self._singleton_initialized = False @classmethod async def reset(cls) -> None: """Reset the singleton instance (mainly for testing).""" with cls._lock: if cls._instance is not None: await cls._instance.close() cls._instance = None # ============= Session methods ============= def session(self, session_id: Optional[str] = None, must_exist: bool = False) -> Session: """ Create a new session or load an existing one. Args: session_id: Session ID, creates a new session (auto-generated ID) if None must_exist: If True and session_id is provided, raises NotFoundError when the session does not exist. If session_id is None, must_exist is ignored. """ return self._client.session(session_id, must_exist=must_exist) async def session_exists(self, session_id: str) -> bool: """Check whether a session exists in storage. Args: session_id: Session ID to check Returns: True if the session exists, False otherwise """ await self._ensure_initialized() return await self._client.session_exists(session_id) async def create_session( self, session_id: Optional[str] = None, telemetry: TelemetryRequest = False, memory_policy: Optional[Dict[str, Any]] = None, ) -> Dict[str, Any]: """Create a new session. Args: session_id: Optional session ID. If provided, creates a session with the given ID. If None, creates a new session with auto-generated ID. """ await self._ensure_initialized() return await self._client.create_session( session_id, telemetry=telemetry, memory_policy=memory_policy, ) async def list_sessions(self) -> List[Any]: """List all sessions.""" await self._ensure_initialized() return await self._client.list_sessions() async def get_session(self, session_id: str, *, auto_create: bool = False) -> Dict[str, Any]: """Get session details.""" await self._ensure_initialized() return await self._client.get_session(session_id, auto_create=auto_create) async def get_session_context( self, session_id: str, token_budget: int = 128_000 ) -> Dict[str, Any]: """Get assembled session context.""" await self._ensure_initialized() return await self._client.get_session_context(session_id, token_budget=token_budget) async def get_session_archive(self, session_id: str, archive_id: str) -> Dict[str, Any]: """Get one completed archive for a session.""" await self._ensure_initialized() return await self._client.get_session_archive(session_id, archive_id) async def delete_session(self, session_id: str) -> None: """Delete a session.""" await self._ensure_initialized() await self._client.delete_session(session_id) async def add_message( self, session_id: str, role: str, content: str | None = None, parts: list[dict] | None = None, created_at: str | None = None, peer_id: str | None = None, telemetry: TelemetryRequest = False, turn_id: str | None = None, message_kind: str | None = None, source_message_ids: list[str] | None = None, ) -> Dict[str, Any]: """Add a message to a session. Args: session_id: Session ID role: Message role ("user" or "assistant") content: Text content (simple mode) parts: Parts array (full Part support: TextPart, ContextPart, ImagePart, ToolPart) created_at: Message creation time (ISO format string) peer_id: Optional stable interaction peer identity. If both content and parts are provided, parts takes precedence. """ await self._ensure_initialized() semantic_kwargs = { key: value for key, value in { "turn_id": turn_id, "message_kind": message_kind, "source_message_ids": source_message_ids, }.items() if value is not None } return await self._client.add_message( session_id=session_id, role=role, content=content, parts=parts, created_at=created_at, peer_id=peer_id, telemetry=telemetry, **semantic_kwargs, ) async def batch_add_messages( self, session_id: str, messages: list[dict], telemetry: TelemetryRequest = False, ) -> Dict[str, Any]: """Add multiple messages to a session in a single request.""" await self._ensure_initialized() return await self._client.batch_add_messages( session_id=session_id, messages=messages, telemetry=telemetry, ) async def commit_session( self, session_id: str, telemetry: TelemetryRequest = False, *, keep_recent_count: int = 0, retention_mode: str | None = None, keep_recent_turn_count: int | None = None, retained_message_token_budget: int | None = None, min_raw_tail_steps: int | None = None, ) -> Dict[str, Any]: """Commit a session (archive and extract memories).""" await self._ensure_initialized() optional_retention = { key: value for key, value in { "retention_mode": retention_mode, "keep_recent_turn_count": keep_recent_turn_count, "retained_message_token_budget": retained_message_token_budget, "min_raw_tail_steps": min_raw_tail_steps, }.items() if value is not None } return await self._client.commit_session( session_id, telemetry=telemetry, keep_recent_count=keep_recent_count, **optional_retention, ) async def get_task(self, task_id: str) -> Optional[Dict[str, Any]]: """Query background task status.""" await self._ensure_initialized() return await self._client.get_task(task_id) async def cancel_task(self, task_id: str) -> Optional[Dict[str, Any]]: """Cancel a background task.""" await self._ensure_initialized() return await self._client.cancel_task(task_id) async def list_tasks( self, task_type: Optional[str] = None, status: Optional[str] = None, resource_id: Optional[str] = None, limit: int = 50, ) -> list[dict[str, Any]]: """List background tasks visible to the current caller.""" await self._ensure_initialized() return await self._client.list_tasks( task_type=task_type, status=status, resource_id=resource_id, limit=limit, ) async def reindex( self, uri: str, mode: str = "vectors_only", wait: bool = True, dry_run: bool = False, ) -> Dict[str, Any]: """Reindex semantic/vector artifacts for a URI.""" await self._ensure_initialized() return await self._client.reindex( uri=uri, mode=mode, wait=wait, dry_run=dry_run, ) # ============= Resource methods ============= async def add_resource( self, path: str, to: Optional[str] = None, parent: Optional[str] = None, reason: str = "", instruction: str = "", wait: bool = False, timeout: float = None, build_index: bool = True, summarize: bool = False, watch_interval: float = 0, 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, ) -> Dict[str, Any]: """ Add a resource (file/URL) to OpenViking. Args: 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. to: Exact target URI. Existing targets keep the add_resource incremental-update behavior. parent: Target parent URI (must already exist). build_index: Whether to build vector index immediately (default: True). summarize: Whether to generate summary (default: False). processing_mode: "semantic_and_vectors" for normal semantic processing, or "vectors_only" to only build vector indexes. watch_interval: Auto-refresh interval in minutes (>0 enables a watch). On a sitemap/feed URL this keeps the whole site refreshed. args: Parser/accessor-specific options (e.g. ``site``, ``max_pages``). telemetry: Whether to attach operation telemetry data to the result. """ 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, instruction=instruction, wait=wait, timeout=timeout, build_index=build_index, summarize=summarize, processing_mode=processing_mode, telemetry=telemetry, watch_interval=watch_interval, args=args, tags=tags, tag_mode=tag_mode, **kwargs, ) @property def _service(self): return self._client.service @property def snapshot(self) -> "AsyncSnapshotNamespace": """Snapshot version control namespace. Lazy-initialized on first access so importing the client does not pull in the snapshot module when it's not needed. """ if getattr(self, "_snapshot", None) is None: from openviking.snapshot_namespace import AsyncSnapshotNamespace self._snapshot = AsyncSnapshotNamespace(self) return self._snapshot async def wait_processed(self, timeout: float = None) -> Dict[str, Any]: """Wait for all queued processing to complete.""" await self._ensure_initialized() return await self._client.wait_processed(timeout=timeout) async def build_index(self, resource_uris: Union[str, List[str]], **kwargs) -> Dict[str, Any]: """ Manually trigger index building for resources. Args: resource_uris: Single URI or list of URIs to index. """ await self._ensure_initialized() return await self._client.build_index(resource_uris, **kwargs) async def summarize(self, resource_uris: Union[str, List[str]], **kwargs) -> Dict[str, Any]: """ Manually trigger summarization for resources. Args: resource_uris: Single URI or list of URIs to summarize. """ await self._ensure_initialized() return await self._client.summarize(resource_uris, **kwargs) async def add_skill( self, data: Any, wait: bool = False, timeout: float = None, telemetry: TelemetryRequest = False, target_uri: Optional[str] = None, ) -> Dict[str, Any]: """Add skill to OpenViking. Args: wait: Whether to wait for vectorization to complete timeout: Wait timeout in seconds target_uri: Optional target root URI override. Defaults to the user's private ``viking://user/{user_id}/skills`` directory. Pass ``viking://agent/skills`` to install a shared skill. """ await self._ensure_initialized() return await self._client.add_skill( data=data, wait=wait, timeout=timeout, telemetry=telemetry, target_uri=target_uri, ) async def list_skills( self, node_limit: int = 1000, target_uri: Optional[str] = None, ) -> Dict[str, Any]: """List installed skills.""" await self._ensure_initialized() return await self._client.list_skills( node_limit=node_limit, target_uri=target_uri, ) async def find_skills( self, query: str, limit: int = 10, score_threshold: Optional[float] = None, level: Optional[List[int]] = None, telemetry: TelemetryRequest = False, target_uri: Optional[str] = None, ) -> Dict[str, Any]: """Find skills by semantic search.""" await self._ensure_initialized() return await self._client.find_skills( query=query, limit=limit, score_threshold=score_threshold, level=level, telemetry=telemetry, target_uri=target_uri, ) async def get_skill( self, skill_name: str, include_content: Optional[bool] = None, include_files: bool = True, include_source: bool = False, level: Optional[int] = None, target_uri: Optional[str] = None, ) -> Dict[str, Any]: """Get a skill by name.""" await self._ensure_initialized() return await self._client.get_skill( skill_name=skill_name, include_content=include_content, include_files=include_files, include_source=include_source, level=level, target_uri=target_uri, ) async def update_skill( self, skill_name: str, data: Any, wait: bool = False, timeout: Optional[float] = None, source_metadata: Optional[Dict[str, Any]] = None, telemetry: TelemetryRequest = False, target_uri: Optional[str] = None, ) -> Dict[str, Any]: """Update an existing skill.""" await self._ensure_initialized() return await self._client.update_skill( skill_name=skill_name, data=data, wait=wait, timeout=timeout, source_metadata=source_metadata, telemetry=telemetry, target_uri=target_uri, ) async def delete_skill( self, skill_name: str, target_uri: Optional[str] = None, ) -> Dict[str, Any]: """Delete a skill.""" await self._ensure_initialized() return await self._client.delete_skill( skill_name=skill_name, target_uri=target_uri, ) async def validate_skill( self, data: Any, strict: bool = False, source_path: Optional[str] = None, skill_dir_name: Optional[str] = None, target_uri: Optional[str] = None, ) -> Dict[str, Any]: """Validate skill data.""" await self._ensure_initialized() return await self._client.validate_skill( data=data, strict=strict, source_path=source_path, skill_dir_name=skill_dir_name, target_uri=target_uri, ) # ============= Search methods ============= async def search( self, query: str = "", target_uri: Union[str, List[str]] = "", session: Optional[Union["Session", Any]] = None, session_id: Optional[str] = None, limit: int = 10, score_threshold: Optional[float] = None, filter: Optional[Dict] = None, context_type: Optional[SearchContextTypeInput] = None, tags: Optional[List[str]] = None, telemetry: TelemetryRequest = False, since: Optional[str] = None, until: Optional[str] = None, time_field: Optional[str] = None, level: Optional[List[int]] = None, image: Optional[Any] = None, ): """ Complex search with session context. Args: query: Query string target_uri: Target directory URI session: Session object for context session_id: Session ID string (alternative to session object) limit: Max results filter: Metadata filters Returns: FindResult """ await self._ensure_initialized() sid = session_id or (session.session_id if session else None) return await self._client.search( query=query, target_uri=target_uri, session_id=sid, limit=limit, score_threshold=score_threshold, filter=filter, context_type=context_type, tags=tags, telemetry=telemetry, since=since, until=until, time_field=time_field, level=level, image=image, ) async def find( self, query: str = "", target_uri: Union[str, List[str]] = "", limit: int = 10, score_threshold: Optional[float] = None, filter: Optional[Dict] = None, context_type: Optional[SearchContextTypeInput] = None, tags: Optional[List[str]] = None, telemetry: TelemetryRequest = False, since: Optional[str] = None, until: Optional[str] = None, time_field: Optional[str] = None, level: Optional[List[int]] = None, image: Optional[Any] = None, ): """Semantic search""" await self._ensure_initialized() return await self._client.find( query=query, target_uri=target_uri, limit=limit, score_threshold=score_threshold, filter=filter, context_type=context_type, tags=tags, telemetry=telemetry, since=since, until=until, time_field=time_field, level=level, image=image, ) # ============= FS methods ============= async def abstract(self, uri: str) -> str: """Read L0 abstract (.abstract.md)""" await self._ensure_initialized() return await self._client.abstract(uri) async def overview(self, uri: str) -> str: """Read L1 overview (.overview.md)""" await self._ensure_initialized() return await self._client.overview(uri) async def read(self, uri: str, offset: int = 0, limit: int = -1) -> str: """Read file content""" await self._ensure_initialized() return await self._client.read(uri, offset=offset, limit=limit) async def read_raw(self, uri: str, offset: int = 0, limit: int = -1) -> str: """Read raw file content, including hidden MEMORY_FIELDS metadata.""" await self._ensure_initialized() read_raw = getattr(self._client, "read_raw", None) if read_raw is not None: return await read_raw(uri, offset=offset, limit=limit) return await self._client.read(uri, offset=offset, limit=limit) async def write( self, uri: str, content: str, mode: str = "replace", wait: bool = False, timeout: Optional[float] = None, telemetry: TelemetryRequest = False, ) -> Dict[str, Any]: """Write text content to an existing file and refresh semantics/vectors.""" await self._ensure_initialized() return await self._client.write( uri=uri, content=content, mode=mode, wait=wait, timeout=timeout, telemetry=telemetry, ) async def set_tags( self, uri: str, tags: List[str], mode: str = "replace", recursive: bool = False, telemetry: TelemetryRequest = False, ) -> Dict[str, Any]: """Replace explicit retrieval tags for a file or directory.""" await self._ensure_initialized() return await self._client.set_tags( uri=uri, tags=tags, mode=mode, recursive=recursive, telemetry=telemetry, ) async def ls(self, uri: str, **kwargs) -> List[Any]: """ List directory contents. Args: uri: Viking URI simple: Return only relative path list (bool, default: False) recursive: List all subdirectories recursively (bool, default: False) node_limit: Maximum number of entries to return (int, default: 1000) sort_by: Optional sort field, "name" or "mtime" sort_order: Sort direction, "asc" or "desc" """ await self._ensure_initialized() recursive = kwargs.get("recursive", False) simple = kwargs.get("simple", False) output = kwargs.get("output", "original") abs_limit = kwargs.get("abs_limit", 256) show_all_hidden = kwargs.get("show_all_hidden", True) node_limit = kwargs.get("node_limit", 1000) sort_by = kwargs.get("sort_by") sort_order = kwargs.get("sort_order", "asc") return await self._client.ls( uri, recursive=recursive, simple=simple, output=output, abs_limit=abs_limit, show_all_hidden=show_all_hidden, node_limit=node_limit, sort_by=sort_by, sort_order=sort_order, ) async def rm( self, uri: str, recursive: bool = False, wait: bool = False, timeout: Optional[float] = None, ) -> None: """Remove resource""" await self._ensure_initialized() await self._client.rm(uri, recursive=recursive, wait=wait, timeout=timeout) async def grep( self, uri: str, pattern: str, case_insensitive: bool = False, node_limit: Optional[int] = None, exclude_uri: Optional[str] = None, level_limit: int = 5, ) -> Dict: """Content search""" await self._ensure_initialized() return await self._client.grep( uri, pattern, case_insensitive=case_insensitive, node_limit=node_limit, exclude_uri=exclude_uri, level_limit=level_limit, ) async def glob(self, pattern: str, uri: str = "viking://") -> Dict: """File pattern matching""" await self._ensure_initialized() return await self._client.glob(pattern, uri=uri) async def mv(self, from_uri: str, to_uri: str) -> None: """Move resource""" await self._ensure_initialized() await self._client.mv(from_uri, to_uri) async def tree(self, uri: str, **kwargs) -> Dict: """Get directory tree""" await self._ensure_initialized() output = kwargs.get("output", "original") abs_limit = kwargs.get("abs_limit", 128) show_all_hidden = kwargs.get("show_all_hidden", True) node_limit = kwargs.get("node_limit", 1000) return await self._client.tree( uri, output=output, abs_limit=abs_limit, show_all_hidden=show_all_hidden, node_limit=node_limit, ) async def mkdir(self, uri: str, description: Optional[str] = None) -> None: """Create directory""" await self._ensure_initialized() await self._client.mkdir(uri, description=description) async def stat(self, uri: str) -> Dict: """Get resource status""" await self._ensure_initialized() return await self._client.stat(uri) # ============= Relation methods ============= async def relations(self, uri: str) -> List[Dict[str, Any]]: """Get relations (returns [{"uri": "...", "reason": "..."}, ...])""" await self._ensure_initialized() return await self._client.relations(uri) async def link(self, from_uri: str, uris: Any, reason: str = "") -> None: """ Create link (single or multiple). Args: from_uri: Source URI uris: Target URI or list of URIs reason: Reason for linking """ await self._ensure_initialized() await self._client.link(from_uri, uris, reason) async def unlink(self, from_uri: str, uri: str) -> None: """ Remove link (remove specified URI from uris). Args: from_uri: Source URI uri: Target URI to remove """ await self._ensure_initialized() await self._client.unlink(from_uri, uri) # ============= Pack methods ============= async def export_ovpack( self, uri: str, to: str, include_vectors: bool = False, ) -> str: """ Export specified context path as .ovpack file. Args: uri: Viking URI to: Target file path Returns: Exported file path """ await self._ensure_initialized() return await self._client.export_ovpack( uri, to, include_vectors=include_vectors, ) async def backup_ovpack(self, to: str, include_vectors: bool = False) -> str: """ Back up public OpenViking scopes as a restore-only .ovpack file. Args: to: Target file path Returns: Exported backup file path """ await self._ensure_initialized() return await self._client.backup_ovpack(to, include_vectors=include_vectors) async def import_ovpack( self, file_path: str, parent: str, on_conflict: Optional[str] = None, vector_mode: Optional[str] = None, ) -> str: """ Import local .ovpack file to specified parent path. Args: file_path: Local .ovpack file path parent: Target parent URI (e.g., viking://user/alice/resources/references/) on_conflict: One of "fail", "overwrite", or "skip" vector_mode: One of "auto", "recompute", or "require" Returns: Imported root resource URI """ await self._ensure_initialized() return await self._client.import_ovpack( file_path, parent, on_conflict=on_conflict, vector_mode=vector_mode, ) async def restore_ovpack( self, file_path: str, on_conflict: Optional[str] = None, vector_mode: Optional[str] = None, ) -> str: """ Restore a backup .ovpack file to its original public scope roots. Args: file_path: Local backup .ovpack file path on_conflict: One of "fail", "overwrite", or "skip" vector_mode: One of "auto", "recompute", or "require" Returns: Restored root URI """ await self._ensure_initialized() return await self._client.restore_ovpack( file_path, on_conflict=on_conflict, vector_mode=vector_mode, ) # ============= Debug methods ============= async def check_consistency(self, uri: str) -> Dict[str, Any]: """Check filesystem/vector-index consistency for a URI subtree.""" await self._ensure_initialized() return await self._client.check_consistency(uri) def get_status(self) -> Union[SystemStatus, Dict[str, Any]]: """Get system status. Returns: SystemStatus containing health status of all components. """ return self._client.get_status() def is_healthy(self) -> bool: """Quick health check. Returns: True if all components are healthy, False otherwise. """ return self._client.is_healthy() @property def observer(self): """Get observer service for component status.""" return self._client.observer