mirror of
https://github.com/volcengine/OpenViking.git
synced 2026-10-01 09:48:03 +08:00
* feat(tasks): support task cancellation * refactor(tasks): scope cancellation to current user * feat(cli): support task cancellation * refactor(tasks): make cancellation queue-aware * refactor(tasks): simplify cancellation bookkeeping * test: remove task cancellation coverage * refactor(tasks): trim cancellation coordination * fix(tasks): contain cancellation to owned work * feat(tasks): persist resource source metadata * fix(tasks): handle cancelled work consistently * refactor(tasks): make completion queue-aware * fix(tasks): persist terminal state before queue ack * test(tasks): remove added lifecycle tests * docs(tasks): document task cancellation
960 lines
32 KiB
Python
960 lines
32 KiB
Python
# 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
|