mirror of
https://github.com/volcengine/OpenViking.git
synced 2026-10-01 09:48:03 +08:00
Keep sessions user-scoped and remove legacy agent fallback that could leak an actor view into commit memory writes.
1363 lines
47 KiB
Python
1363 lines
47 KiB
Python
# Copyright (c) 2026 Beijing Volcano Engine Technology Co., Ltd.
|
|
# SPDX-License-Identifier: AGPL-3.0
|
|
"""Local Client for OpenViking.
|
|
|
|
Implements BaseClient interface using direct service calls (embedded mode).
|
|
"""
|
|
|
|
from typing import Any, Dict, List, Optional, Union
|
|
|
|
from openviking.core.peer_id import normalize_peer_id
|
|
from openviking.core.skill_loader import validate_skill_format
|
|
from openviking.server.identity import RequestContext, Role
|
|
from openviking.server.routers.skills import (
|
|
_list_skill_files,
|
|
_list_skills_from_root,
|
|
_require_skill,
|
|
_restore_skill_privacy,
|
|
_skill_summary_from_hit,
|
|
)
|
|
from openviking.service import OpenVikingService
|
|
from openviking.service.task_tracker import get_task_tracker
|
|
from openviking.telemetry import TelemetryRequest
|
|
from openviking.telemetry.execution import (
|
|
attach_telemetry_payload,
|
|
run_with_telemetry,
|
|
)
|
|
from openviking.utils.image_search import normalize_client_image_input
|
|
from openviking.utils.search_filters import SearchContextTypeInput, merge_search_filter
|
|
from openviking.utils.tags import build_search_tags_filter
|
|
from openviking_cli.client.base import BaseClient
|
|
from openviking_cli.exceptions import InvalidArgumentError, NotFoundError, PermissionDeniedError
|
|
from openviking_cli.session.user_id import UserIdentifier
|
|
from openviking_cli.utils import run_async
|
|
|
|
|
|
def _to_jsonable(value: Any) -> Any:
|
|
"""Convert internal objects into JSON-serializable values."""
|
|
to_dict = getattr(value, "to_dict", None)
|
|
if callable(to_dict):
|
|
return to_dict()
|
|
if isinstance(value, list):
|
|
return [_to_jsonable(item) for item in value]
|
|
if isinstance(value, dict):
|
|
return {k: _to_jsonable(v) for k, v in value.items()}
|
|
return value
|
|
|
|
|
|
def _resolve_search_filter(
|
|
filter: Optional[Dict[str, Any]],
|
|
context_type: Optional[SearchContextTypeInput],
|
|
since: Optional[str],
|
|
until: Optional[str],
|
|
time_field: Optional[str],
|
|
tags: Optional[List[str]] = None,
|
|
) -> Optional[Dict[str, Any]]:
|
|
"""Merge public retrieval filter shortcuts into the metadata filter."""
|
|
merged = merge_search_filter(
|
|
filter,
|
|
context_type=context_type,
|
|
since=since,
|
|
until=until,
|
|
time_field=time_field,
|
|
)
|
|
tag_filter = build_search_tags_filter(tags)
|
|
if not tag_filter:
|
|
return merged
|
|
if merged:
|
|
if tag_filter.get("op") == "and" and isinstance(tag_filter.get("conds"), list):
|
|
return {"op": "and", "conds": [merged, *tag_filter["conds"]]}
|
|
return {"op": "and", "conds": [merged, tag_filter]}
|
|
return tag_filter
|
|
|
|
|
|
class LocalClient(BaseClient):
|
|
"""Local Client for OpenViking (embedded mode).
|
|
|
|
Implements BaseClient interface using direct service calls.
|
|
"""
|
|
|
|
def __init__(
|
|
self,
|
|
path: Optional[str] = None,
|
|
user: Optional[UserIdentifier] = None,
|
|
actor_peer_id: Optional[str] = None,
|
|
agent_id: Optional[str] = None,
|
|
):
|
|
"""Initialize LocalClient.
|
|
|
|
Args:
|
|
path: Local storage path (overrides ov.conf storage path)
|
|
user: Explicit account/user identity for embedded mode
|
|
actor_peer_id: Optional view filter for the current user's peer collection.
|
|
agent_id: Legacy alias that marks the actor peer scope as legacy agent mode.
|
|
"""
|
|
if actor_peer_id is not None and agent_id is not None:
|
|
raise ValueError("actor_peer_id cannot be used with legacy agent_id")
|
|
effective_actor_peer_id = actor_peer_id or agent_id
|
|
self._service = OpenVikingService(
|
|
path=path,
|
|
user=user or UserIdentifier.the_default_user(),
|
|
)
|
|
self._user = self._service.user
|
|
self._ctx = RequestContext(
|
|
user=self._user,
|
|
role=Role.USER,
|
|
actor_peer_id=normalize_peer_id(effective_actor_peer_id),
|
|
)
|
|
|
|
@property
|
|
def service(self) -> OpenVikingService:
|
|
"""Get the underlying service instance."""
|
|
return self._service
|
|
|
|
# ============= Lifecycle =============
|
|
|
|
async def initialize(self) -> None:
|
|
"""Initialize the local client."""
|
|
await self._service.initialize()
|
|
|
|
async def close(self) -> None:
|
|
"""Close the local client."""
|
|
await self._service.close()
|
|
|
|
# ============= Resource Management =============
|
|
|
|
async def add_resource(
|
|
self,
|
|
path: str,
|
|
to: Optional[str] = None,
|
|
parent: Optional[str] = None,
|
|
reason: str = "",
|
|
instruction: str = "",
|
|
wait: bool = False,
|
|
timeout: Optional[float] = None,
|
|
build_index: bool = True,
|
|
summarize: bool = False,
|
|
telemetry: TelemetryRequest = False,
|
|
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_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.")
|
|
|
|
execution = await run_with_telemetry(
|
|
operation="resources.add_resource",
|
|
telemetry=telemetry,
|
|
fn=lambda: self._service.resources.add_resource(
|
|
path=path,
|
|
ctx=self._ctx,
|
|
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,
|
|
watch_interval=watch_interval,
|
|
args=args,
|
|
tags=tags,
|
|
tag_mode=tag_mode,
|
|
**kwargs,
|
|
),
|
|
)
|
|
return attach_telemetry_payload(
|
|
execution.result,
|
|
execution.telemetry,
|
|
)
|
|
|
|
async def add_skill(
|
|
self,
|
|
data: Any,
|
|
wait: bool = False,
|
|
timeout: Optional[float] = None,
|
|
telemetry: TelemetryRequest = False,
|
|
target_uri: Optional[str] = None,
|
|
) -> Dict[str, Any]:
|
|
"""Add skill to OpenViking."""
|
|
execution = await run_with_telemetry(
|
|
operation="resources.add_skill",
|
|
telemetry=telemetry,
|
|
fn=lambda: self._service.resources.add_skill(
|
|
data=data,
|
|
ctx=self._ctx,
|
|
wait=wait,
|
|
timeout=timeout,
|
|
target_uri=target_uri,
|
|
),
|
|
)
|
|
return attach_telemetry_payload(
|
|
execution.result,
|
|
execution.telemetry,
|
|
)
|
|
|
|
async def list_skills(
|
|
self,
|
|
node_limit: int = 1000,
|
|
target_uri: Optional[str] = None,
|
|
) -> Dict[str, Any]:
|
|
"""List installed skills."""
|
|
from openviking.core.namespace import canonical_user_root
|
|
|
|
service = self._service
|
|
ctx = self._ctx
|
|
if target_uri:
|
|
skills = await _list_skills_from_root(service, ctx, target_uri)
|
|
return {"root_uri": target_uri, "skills": skills, "total": len(skills)}
|
|
else:
|
|
user_skills = await _list_skills_from_root(
|
|
service, ctx, f"{canonical_user_root(ctx)}/skills"
|
|
)
|
|
agent_skills = await _list_skills_from_root(service, ctx, "viking://agent/skills")
|
|
merged = [*user_skills, *agent_skills]
|
|
return {
|
|
"root_uris": [f"{canonical_user_root(ctx)}/skills", "viking://agent/skills"],
|
|
"skills": merged,
|
|
"total": len(merged),
|
|
}
|
|
|
|
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."""
|
|
from openviking.core.namespace import canonical_user_root
|
|
|
|
service = self._service
|
|
ctx = self._ctx
|
|
|
|
async def _search_at(uri: str) -> list:
|
|
result = await service.search.find(
|
|
query=query,
|
|
ctx=ctx,
|
|
target_uri=uri,
|
|
limit=limit,
|
|
score_threshold=score_threshold,
|
|
level=level,
|
|
)
|
|
result_dict = result.to_dict() if hasattr(result, "to_dict") else dict(result or {})
|
|
return [_skill_summary_from_hit(hit) for hit in result_dict.get("skills", [])]
|
|
|
|
if target_uri:
|
|
execution = await run_with_telemetry(
|
|
operation="skills.find",
|
|
telemetry=telemetry,
|
|
fn=lambda: _search_at(target_uri),
|
|
)
|
|
hits = execution.result
|
|
return {
|
|
"root_uri": target_uri,
|
|
"skills": hits,
|
|
"total": len(hits),
|
|
}
|
|
else:
|
|
user_root = f"{canonical_user_root(ctx)}/skills"
|
|
agent_root = "viking://agent/skills"
|
|
|
|
user_execution = await run_with_telemetry(
|
|
operation="skills.find",
|
|
telemetry=telemetry,
|
|
fn=lambda: _search_at(user_root),
|
|
)
|
|
user_hits = user_execution.result
|
|
|
|
agent_hits = await _search_at(agent_root)
|
|
|
|
merged = [*user_hits, *agent_hits]
|
|
if merged and "score" in merged[0]:
|
|
merged.sort(key=lambda x: x.get("score", 0), reverse=True)
|
|
|
|
return {
|
|
"root_uris": [user_root, agent_root],
|
|
"skills": merged,
|
|
"total": len(merged),
|
|
}
|
|
|
|
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."""
|
|
from openviking.server.routers.skills import (
|
|
SOURCE_METADATA_FILENAME,
|
|
_parse_abstract_meta,
|
|
_relative_skill_path,
|
|
_skill_file_kind,
|
|
_skill_summary_from_meta,
|
|
)
|
|
from openviking.server.skill_source_metadata import read_skill_source_metadata
|
|
|
|
if level is not None and level not in {0, 1, 2}:
|
|
raise InvalidArgumentError(
|
|
"Skill show level must be 0, 1, or 2",
|
|
details={"field": "level", "allowed": [0, 1, 2]},
|
|
)
|
|
|
|
service = self._service
|
|
ctx = self._ctx
|
|
root_uri = await _require_skill(service, ctx, skill_name, target_uri)
|
|
|
|
abstract = await service.fs.abstract(root_uri, ctx=ctx)
|
|
result = _skill_summary_from_meta(skill_name, root_uri, _parse_abstract_meta(abstract))
|
|
|
|
if level is None or level == 0:
|
|
result["abstract"] = abstract
|
|
if level is None or level == 1:
|
|
result["overview"] = await service.fs.overview(root_uri, ctx=ctx)
|
|
if (
|
|
level == 2
|
|
or include_content is True
|
|
or (level is None and include_content is not False)
|
|
):
|
|
from openviking.server.routers.skills import _skill_md_uri
|
|
|
|
result["content"] = await service.fs.read(_skill_md_uri(root_uri), ctx=ctx)
|
|
|
|
if include_files:
|
|
entries = await _list_skill_files(service, ctx, root_uri)
|
|
result["files"] = [
|
|
{
|
|
"name": entry.get("name") or skill_name,
|
|
"uri": entry.get("uri", ""),
|
|
"path": _relative_skill_path(root_uri, entry.get("uri", "")),
|
|
"is_dir": entry.get("isDir", False),
|
|
"kind": _skill_file_kind(
|
|
_relative_skill_path(root_uri, entry.get("uri", "")),
|
|
entry.get("isDir", False),
|
|
),
|
|
}
|
|
for entry in entries
|
|
if isinstance(entry, dict)
|
|
and _relative_skill_path(root_uri, entry.get("uri", "")) != SOURCE_METADATA_FILENAME
|
|
]
|
|
|
|
if include_source:
|
|
result["source"] = await read_skill_source_metadata(service, ctx, root_uri)
|
|
|
|
return result
|
|
|
|
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."""
|
|
service = self._service
|
|
ctx = self._ctx
|
|
|
|
# Verify the skill exists and determine its root URI
|
|
root_uri = await _require_skill(service, ctx, skill_name, target_uri)
|
|
skill_root_parent = root_uri.rsplit("/", 1)[0]
|
|
|
|
execution = await run_with_telemetry(
|
|
operation="skills.update",
|
|
telemetry=telemetry,
|
|
fn=lambda: self._update_skill_impl(
|
|
skill_name, data, root_uri, skill_root_parent, wait, timeout, source_metadata
|
|
),
|
|
)
|
|
return attach_telemetry_payload(
|
|
execution.result,
|
|
execution.telemetry,
|
|
)
|
|
|
|
async def _update_skill_impl(
|
|
self,
|
|
skill_name: str,
|
|
data: Any,
|
|
root_uri: str,
|
|
skill_root_parent: str,
|
|
wait: bool,
|
|
timeout: Optional[float],
|
|
source_metadata: Optional[Dict[str, Any]],
|
|
) -> Dict[str, Any]:
|
|
import shutil
|
|
import uuid
|
|
|
|
from openviking.server.skill_source_metadata import persist_skill_source_metadata
|
|
|
|
service = self._service
|
|
ctx = self._ctx
|
|
backup_uri = f"{skill_root_parent}/.{skill_name}.update-backup-{uuid.uuid4().hex}"
|
|
backup_created = False
|
|
privacy_update_attempted = False
|
|
previous_privacy = None
|
|
preparation = None
|
|
privacy = service.privacy_configs
|
|
effective_source_metadata = source_metadata or {
|
|
"type": "embedded",
|
|
"source": "inline_content",
|
|
"operation": "update",
|
|
}
|
|
try:
|
|
if privacy is not None:
|
|
previous_privacy = await privacy.get_current(ctx, "skill", skill_name)
|
|
preparation = await service.resources._skill_processor.prepare_skill_processing( # noqa: SLF001
|
|
data,
|
|
ctx=ctx,
|
|
allow_local_path_resolution=False,
|
|
)
|
|
expected_name = skill_name
|
|
if preparation.skill_dict.get("name") != expected_name:
|
|
raise InvalidArgumentError(
|
|
f"Skill name mismatch: path name is '{expected_name}', content name is '{preparation.skill_dict.get('name')}'",
|
|
details={
|
|
"expected": expected_name,
|
|
"actual": preparation.skill_dict.get("name"),
|
|
},
|
|
)
|
|
await service.fs.mv(root_uri, backup_uri, ctx=ctx)
|
|
backup_created = True
|
|
result = await service.resources.add_skill(
|
|
data=preparation,
|
|
ctx=ctx,
|
|
wait=wait,
|
|
timeout=timeout,
|
|
allow_local_path_resolution=False,
|
|
apply_privacy=False,
|
|
privacy_change_reason="auto-extracted from update_skill",
|
|
target_uri=skill_root_parent,
|
|
)
|
|
await persist_skill_source_metadata(service, ctx, result, effective_source_metadata)
|
|
privacy_update_attempted = True
|
|
await service.resources._skill_processor.apply_skill_privacy( # noqa: SLF001
|
|
preparation.skill_dict,
|
|
preparation.privacy_values,
|
|
ctx,
|
|
change_reason="auto-extracted from update_skill",
|
|
delete_if_empty=True,
|
|
)
|
|
except Exception:
|
|
if backup_created:
|
|
try:
|
|
await service.fs.rm(root_uri, ctx=ctx, recursive=True)
|
|
except Exception:
|
|
pass
|
|
try:
|
|
await service.fs.mv(backup_uri, root_uri, ctx=ctx)
|
|
except Exception:
|
|
pass
|
|
if privacy_update_attempted:
|
|
try:
|
|
await _restore_skill_privacy(service, ctx, skill_name, previous_privacy)
|
|
except Exception:
|
|
pass
|
|
raise
|
|
else:
|
|
if backup_created:
|
|
try:
|
|
await service.fs.rm(backup_uri, ctx=ctx, recursive=True)
|
|
except Exception:
|
|
pass
|
|
result["action"] = "update"
|
|
return result
|
|
finally:
|
|
if preparation and preparation.cleanup_path:
|
|
shutil.rmtree(preparation.cleanup_path, ignore_errors=True)
|
|
|
|
async def delete_skill(
|
|
self,
|
|
skill_name: str,
|
|
target_uri: Optional[str] = None,
|
|
) -> Dict[str, Any]:
|
|
"""Delete a skill."""
|
|
service = self._service
|
|
ctx = self._ctx
|
|
root_uri = await _require_skill(service, ctx, skill_name, target_uri)
|
|
await service.fs.rm(root_uri, ctx=ctx, recursive=True)
|
|
return {"name": skill_name, "uri": root_uri, "root_uri": root_uri, "deleted": True}
|
|
|
|
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."""
|
|
result = validate_skill_format(
|
|
data,
|
|
strict=strict,
|
|
skill_dir_name=skill_dir_name,
|
|
source_path=source_path,
|
|
)
|
|
return result
|
|
|
|
async def wait_processed(self, timeout: Optional[float] = None) -> Dict[str, Any]:
|
|
"""Wait for all processing to complete."""
|
|
return await self._service.resources.wait_processed(timeout=timeout)
|
|
|
|
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."""
|
|
return await self._service.reindex(
|
|
uri=uri,
|
|
mode=mode,
|
|
wait=wait,
|
|
dry_run=dry_run,
|
|
)
|
|
|
|
async def build_index(self, resource_uris: Union[str, List[str]], **kwargs) -> Dict[str, Any]:
|
|
"""Manually trigger index building."""
|
|
if isinstance(resource_uris, str):
|
|
resource_uris = [resource_uris]
|
|
return await self._service.resources.build_index(resource_uris, ctx=self._ctx, **kwargs)
|
|
|
|
async def summarize(self, resource_uris: Union[str, List[str]], **kwargs) -> Dict[str, Any]:
|
|
"""Manually trigger summarization."""
|
|
if isinstance(resource_uris, str):
|
|
resource_uris = [resource_uris]
|
|
return await self._service.resources.summarize(resource_uris, ctx=self._ctx, **kwargs)
|
|
|
|
# ============= File System =============
|
|
|
|
async def ls(
|
|
self,
|
|
uri: str,
|
|
simple: bool = False,
|
|
recursive: bool = False,
|
|
output: str = "original",
|
|
abs_limit: int = 256,
|
|
show_all_hidden: bool = False,
|
|
node_limit: int = 1000,
|
|
sort_by: Optional[str] = None,
|
|
sort_order: str = "asc",
|
|
) -> List[Any]:
|
|
"""List directory contents."""
|
|
return await self._service.fs.ls(
|
|
uri,
|
|
ctx=self._ctx,
|
|
simple=simple,
|
|
recursive=recursive,
|
|
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 tree(
|
|
self,
|
|
uri: str,
|
|
output: str = "original",
|
|
abs_limit: int = 128,
|
|
show_all_hidden: bool = False,
|
|
node_limit: int = 1000,
|
|
) -> List[Dict[str, Any]]:
|
|
"""Get directory tree."""
|
|
return await self._service.fs.tree(
|
|
uri,
|
|
ctx=self._ctx,
|
|
output=output,
|
|
abs_limit=abs_limit,
|
|
show_all_hidden=show_all_hidden,
|
|
node_limit=node_limit,
|
|
)
|
|
|
|
async def stat(self, uri: str) -> Dict[str, Any]:
|
|
"""Get resource status."""
|
|
return await self._service.fs.stat(uri, ctx=self._ctx)
|
|
|
|
async def mkdir(self, uri: str, description: Optional[str] = None) -> None:
|
|
"""Create directory."""
|
|
await self._service.fs.mkdir(uri, ctx=self._ctx, description=description)
|
|
|
|
async def rm(
|
|
self,
|
|
uri: str,
|
|
recursive: bool = False,
|
|
wait: bool = False,
|
|
timeout: Optional[float] = None,
|
|
) -> None:
|
|
"""Remove resource."""
|
|
await self._service.fs.rm(
|
|
uri,
|
|
ctx=self._ctx,
|
|
recursive=recursive,
|
|
wait=wait,
|
|
timeout=timeout,
|
|
)
|
|
|
|
async def mv(self, from_uri: str, to_uri: str) -> None:
|
|
"""Move resource."""
|
|
await self._service.fs.mv(from_uri, to_uri, ctx=self._ctx)
|
|
|
|
# ============= Content Reading =============
|
|
|
|
async def read(self, uri: str, offset: int = 0, limit: int = -1) -> str:
|
|
"""Read file content.
|
|
|
|
Args:
|
|
uri: Viking URI
|
|
offset: Starting line number (0-indexed). Default 0.
|
|
limit: Number of lines to read. -1 means read to end. Default -1.
|
|
"""
|
|
return await self._service.fs.read(uri, ctx=self._ctx, 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."""
|
|
return await self._service.fs.read(uri, ctx=self._ctx, offset=offset, limit=limit)
|
|
|
|
async def abstract(self, uri: str) -> str:
|
|
"""Read L0 abstract."""
|
|
return await self._service.fs.abstract(uri, ctx=self._ctx)
|
|
|
|
async def overview(self, uri: str) -> str:
|
|
"""Read L1 overview."""
|
|
return await self._service.fs.overview(uri, ctx=self._ctx)
|
|
|
|
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."""
|
|
execution = await run_with_telemetry(
|
|
operation="content.write",
|
|
telemetry=telemetry,
|
|
fn=lambda: self._service.fs.write(
|
|
uri=uri,
|
|
content=content,
|
|
ctx=self._ctx,
|
|
mode=mode,
|
|
wait=wait,
|
|
timeout=timeout,
|
|
),
|
|
)
|
|
return attach_telemetry_payload(
|
|
execution.result,
|
|
execution.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."""
|
|
execution = await run_with_telemetry(
|
|
operation="content.set_tags",
|
|
telemetry=telemetry,
|
|
fn=lambda: self._service.fs.set_tags(
|
|
uri=uri,
|
|
tags=tags,
|
|
mode=mode,
|
|
recursive=recursive,
|
|
ctx=self._ctx,
|
|
),
|
|
)
|
|
return attach_telemetry_payload(
|
|
execution.result,
|
|
execution.telemetry,
|
|
)
|
|
|
|
# ============= Search =============
|
|
|
|
async def find(
|
|
self,
|
|
query: str = "",
|
|
target_uri: Union[str, List[str]] = "",
|
|
limit: int = 10,
|
|
score_threshold: Optional[float] = None,
|
|
filter: Optional[Dict[str, Any]] = 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,
|
|
) -> Any:
|
|
"""Semantic search without session context."""
|
|
resolved_filter = _resolve_search_filter(
|
|
filter, context_type, since, until, time_field, tags
|
|
)
|
|
image_url = normalize_client_image_input(image)
|
|
execution = await run_with_telemetry(
|
|
operation="search.find",
|
|
telemetry=telemetry,
|
|
fn=lambda: self._service.search.find(
|
|
query=query,
|
|
ctx=self._ctx,
|
|
target_uri=target_uri,
|
|
limit=limit,
|
|
score_threshold=score_threshold,
|
|
filter=resolved_filter,
|
|
level=level,
|
|
image_url=image_url,
|
|
),
|
|
)
|
|
return attach_telemetry_payload(
|
|
execution.result,
|
|
execution.telemetry,
|
|
)
|
|
|
|
async def search(
|
|
self,
|
|
query: str = "",
|
|
target_uri: Union[str, List[str]] = "",
|
|
session_id: Optional[str] = None,
|
|
limit: int = 10,
|
|
score_threshold: Optional[float] = None,
|
|
filter: Optional[Dict[str, Any]] = 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,
|
|
) -> Any:
|
|
"""Semantic search with optional session context."""
|
|
resolved_filter = _resolve_search_filter(
|
|
filter, context_type, since, until, time_field, tags
|
|
)
|
|
image_url = normalize_client_image_input(image)
|
|
|
|
async def _search():
|
|
session = None
|
|
# Intent off: skip session.load — SearchService will not scan session either.
|
|
if session_id and self._service.search.is_intent_enabled():
|
|
session = self._service.sessions.session(self._ctx, session_id)
|
|
await session.load()
|
|
return await self._service.search.search(
|
|
query=query,
|
|
ctx=self._ctx,
|
|
target_uri=target_uri,
|
|
session=session,
|
|
limit=limit,
|
|
score_threshold=score_threshold,
|
|
filter=resolved_filter,
|
|
level=level,
|
|
image_url=image_url,
|
|
)
|
|
|
|
execution = await run_with_telemetry(
|
|
operation="search.search",
|
|
telemetry=telemetry,
|
|
fn=_search,
|
|
)
|
|
return attach_telemetry_payload(
|
|
execution.result,
|
|
execution.telemetry,
|
|
)
|
|
|
|
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[str, Any]:
|
|
"""Content search with pattern."""
|
|
return await self._service.fs.grep(
|
|
uri,
|
|
pattern,
|
|
ctx=self._ctx,
|
|
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[str, Any]:
|
|
"""File pattern matching."""
|
|
return await self._service.fs.glob(pattern, ctx=self._ctx, uri=uri)
|
|
|
|
# ============= Relations =============
|
|
|
|
async def relations(self, uri: str) -> List[Any]:
|
|
"""Get relations for a resource."""
|
|
return await self._service.relations.relations(uri, ctx=self._ctx)
|
|
|
|
async def link(self, from_uri: str, to_uris: Union[str, List[str]], reason: str = "") -> None:
|
|
"""Create link between resources."""
|
|
await self._service.relations.link(from_uri, to_uris, ctx=self._ctx, reason=reason)
|
|
|
|
async def unlink(self, from_uri: str, to_uri: str) -> None:
|
|
"""Remove link between resources."""
|
|
await self._service.relations.unlink(from_uri, to_uri, ctx=self._ctx)
|
|
|
|
# ============= Sessions =============
|
|
|
|
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.
|
|
"""
|
|
execution = await run_with_telemetry(
|
|
operation="session.create",
|
|
telemetry=telemetry,
|
|
fn=lambda: self._create_session_impl(session_id, memory_policy),
|
|
)
|
|
return attach_telemetry_payload(
|
|
execution.result,
|
|
execution.telemetry,
|
|
)
|
|
|
|
async def _create_session_impl(
|
|
self,
|
|
session_id: Optional[str],
|
|
memory_policy: Optional[Dict[str, Any]],
|
|
) -> Dict[str, Any]:
|
|
await self._service.initialize_user_directories(self._ctx)
|
|
session = await self._service.sessions.create(
|
|
self._ctx,
|
|
session_id,
|
|
memory_policy=memory_policy,
|
|
)
|
|
return {
|
|
"session_id": session.session_id,
|
|
"uri": session.uri,
|
|
"user": session.user.to_dict(),
|
|
}
|
|
|
|
async def list_sessions(self) -> List[Any]:
|
|
"""List all sessions."""
|
|
return await self._service.sessions.sessions(self._ctx)
|
|
|
|
async def get_session(self, session_id: str, *, auto_create: bool = False) -> Dict[str, Any]:
|
|
"""Get session details."""
|
|
session = await self._service.sessions.get(session_id, self._ctx, auto_create=auto_create)
|
|
result = session.meta.to_dict()
|
|
result["uri"] = session.uri
|
|
result["user"] = session.user.to_dict()
|
|
return result
|
|
|
|
async def get_session_context(
|
|
self, session_id: str, token_budget: int = 128_000
|
|
) -> Dict[str, Any]:
|
|
"""Get assembled session context."""
|
|
session = await self._service.sessions.get(session_id, self._ctx, auto_create=False)
|
|
result = await session.get_session_context(token_budget=token_budget)
|
|
return _to_jsonable(result)
|
|
|
|
async def get_session_archive(self, session_id: str, archive_id: str) -> Dict[str, Any]:
|
|
"""Get one completed archive for a session."""
|
|
session = await self._service.sessions.get(session_id, self._ctx, auto_create=False)
|
|
result = await session.get_session_archive(archive_id)
|
|
return _to_jsonable(result)
|
|
|
|
async def delete_session(self, session_id: str) -> None:
|
|
"""Delete a session."""
|
|
await self._service.sessions.delete(session_id, self._ctx)
|
|
|
|
async def commit_session(
|
|
self,
|
|
session_id: str,
|
|
telemetry: TelemetryRequest = False,
|
|
*,
|
|
keep_recent_count: int = 0,
|
|
retention_mode: Optional[str] = None,
|
|
keep_recent_turn_count: Optional[int] = None,
|
|
retained_message_token_budget: Optional[int] = None,
|
|
min_raw_tail_steps: Optional[int] = None,
|
|
) -> Dict[str, Any]:
|
|
"""Commit a session (archive and extract memories)."""
|
|
commit_kwargs: Dict[str, Any] = {"keep_recent_count": keep_recent_count}
|
|
optional_retention = {
|
|
"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,
|
|
}
|
|
commit_kwargs.update(
|
|
{key: value for key, value in optional_retention.items() if value is not None}
|
|
)
|
|
execution = await run_with_telemetry(
|
|
operation="session.commit",
|
|
telemetry=telemetry,
|
|
fn=lambda: self._service.sessions.commit(
|
|
session_id,
|
|
self._ctx,
|
|
**commit_kwargs,
|
|
),
|
|
)
|
|
return attach_telemetry_payload(
|
|
execution.result,
|
|
execution.telemetry,
|
|
)
|
|
|
|
async def get_task(self, task_id: str) -> Optional[Dict[str, Any]]:
|
|
"""Query background task status."""
|
|
return await self._service.sessions.get_commit_task(task_id, self._ctx)
|
|
|
|
async def cancel_task(self, task_id: str) -> Optional[Dict[str, Any]]:
|
|
"""Cancel a background task."""
|
|
if self._ctx.role == Role.ROOT:
|
|
raise PermissionDeniedError("ROOT may not cancel tasks")
|
|
task = await get_task_tracker().cancel(
|
|
task_id,
|
|
account_id=self._ctx.account_id,
|
|
user_id=self._ctx.user.user_id,
|
|
)
|
|
return task.to_dict() if task else None
|
|
|
|
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."""
|
|
tasks = await get_task_tracker().list_tasks(
|
|
task_type=task_type,
|
|
status=status,
|
|
resource_id=resource_id,
|
|
limit=limit,
|
|
account_id=self._ctx.account_id,
|
|
user_id=self._ctx.user.user_id,
|
|
)
|
|
return [task.to_dict() for task in tasks]
|
|
|
|
async def add_message(
|
|
self,
|
|
session_id: str,
|
|
role: str,
|
|
content: Optional[str] = None,
|
|
parts: Optional[List[Dict[str, Any]]] = None,
|
|
created_at: Optional[str] = None,
|
|
peer_id: Optional[str] = None,
|
|
telemetry: TelemetryRequest = False,
|
|
turn_id: Optional[str] = None,
|
|
message_kind: Optional[str] = None,
|
|
source_message_ids: Optional[List[str]] = 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, backward compatible)
|
|
parts: Parts array (full Part support mode)
|
|
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.
|
|
"""
|
|
execution = await run_with_telemetry(
|
|
operation="session.add_message",
|
|
telemetry=telemetry,
|
|
fn=lambda: self._add_message_impl(
|
|
session_id,
|
|
role,
|
|
content,
|
|
parts,
|
|
created_at,
|
|
peer_id,
|
|
turn_id,
|
|
message_kind,
|
|
source_message_ids,
|
|
),
|
|
)
|
|
return attach_telemetry_payload(
|
|
execution.result,
|
|
execution.telemetry,
|
|
)
|
|
|
|
async def _add_message_impl(
|
|
self,
|
|
session_id: str,
|
|
role: str,
|
|
content: Optional[str],
|
|
parts: Optional[List[Dict[str, Any]]],
|
|
created_at: Optional[str],
|
|
peer_id: Optional[str],
|
|
turn_id: Optional[str],
|
|
message_kind: Optional[str],
|
|
source_message_ids: Optional[List[str]],
|
|
) -> Dict[str, Any]:
|
|
from openviking.message.part import Part, TextPart, part_from_dict
|
|
|
|
session = await self._service.sessions.get(session_id, self._ctx, auto_create=True)
|
|
|
|
message_parts: list[Part]
|
|
if parts is not None:
|
|
message_parts = [part_from_dict(p) for p in parts]
|
|
elif content is not None:
|
|
message_parts = [TextPart(text=content)]
|
|
else:
|
|
raise ValueError("Either content or parts must be provided")
|
|
|
|
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
|
|
}
|
|
add_async = getattr(session, "add_message_async", None)
|
|
add_kwargs = {
|
|
"peer_id": normalize_peer_id(peer_id),
|
|
"created_at": created_at,
|
|
**semantic_kwargs,
|
|
}
|
|
if callable(add_async):
|
|
await add_async(role, message_parts, **add_kwargs)
|
|
else:
|
|
session.add_message(role, message_parts, **add_kwargs)
|
|
return {
|
|
"session_id": session_id,
|
|
"message_count": len(session.messages),
|
|
# Post-write value so a commit policy can decide without a
|
|
# follow-up get_session round trip.
|
|
"pending_tokens": self._session_pending_tokens(session),
|
|
}
|
|
|
|
async def batch_add_messages(
|
|
self,
|
|
session_id: str,
|
|
messages: List[Dict[str, Any]],
|
|
telemetry: TelemetryRequest = False,
|
|
) -> Dict[str, Any]:
|
|
"""Add multiple messages to a session in one batch."""
|
|
execution = await run_with_telemetry(
|
|
operation="session.batch_add_messages",
|
|
telemetry=telemetry,
|
|
fn=lambda: self._batch_add_messages_impl(session_id, messages),
|
|
)
|
|
return attach_telemetry_payload(
|
|
execution.result,
|
|
execution.telemetry,
|
|
)
|
|
|
|
async def _batch_add_messages_impl(
|
|
self,
|
|
session_id: str,
|
|
messages: List[Dict[str, Any]],
|
|
) -> Dict[str, Any]:
|
|
from openviking.message.part import Part, TextPart, part_from_dict
|
|
|
|
session = await self._service.sessions.get(session_id, self._ctx, auto_create=True)
|
|
specs: list[dict[str, Any]] = []
|
|
|
|
for index, message in enumerate(messages):
|
|
role = message.get("role")
|
|
if not role:
|
|
raise ValueError(f"messages[{index}]: missing required key 'role'")
|
|
|
|
message_parts: list[Part]
|
|
if message.get("parts") is not None:
|
|
message_parts = [part_from_dict(part) for part in message["parts"]]
|
|
elif message.get("content") is not None:
|
|
message_parts = [TextPart(text=str(message["content"]))]
|
|
else:
|
|
raise ValueError(f"messages[{index}]: Either content or parts must be provided")
|
|
|
|
specs.append(
|
|
{
|
|
"role": role,
|
|
"parts": message_parts,
|
|
"peer_id": normalize_peer_id(message.get("peer_id")),
|
|
"created_at": message.get("created_at"),
|
|
"turn_id": message.get("turn_id"),
|
|
"message_kind": message.get("message_kind"),
|
|
"source_message_ids": message.get("source_message_ids"),
|
|
}
|
|
)
|
|
|
|
add_many_async = getattr(session, "add_messages_async", None)
|
|
if callable(add_many_async):
|
|
added = await add_many_async(specs)
|
|
else:
|
|
added = session.add_messages(specs)
|
|
return {
|
|
"session_id": session_id,
|
|
"message_count": len(session.messages),
|
|
"added": len(added),
|
|
# Post-write value so a commit policy can decide without a
|
|
# follow-up get_session round trip.
|
|
"pending_tokens": self._session_pending_tokens(session),
|
|
}
|
|
|
|
@staticmethod
|
|
def _session_pending_tokens(session: Any) -> int:
|
|
"""Read the post-write pending-token count from a session.
|
|
|
|
Returns 0 when the session object does not expose ``meta`` so callers
|
|
keep working against lightweight or legacy session implementations.
|
|
"""
|
|
meta = getattr(session, "meta", None)
|
|
try:
|
|
return max(0, int(getattr(meta, "pending_tokens", 0) or 0))
|
|
except (TypeError, ValueError):
|
|
return 0
|
|
|
|
# ============= Pack =============
|
|
|
|
async def export_ovpack(
|
|
self,
|
|
uri: str,
|
|
to: str,
|
|
include_vectors: bool = False,
|
|
) -> str:
|
|
"""Export context as .ovpack file."""
|
|
return await self._service.pack.export_ovpack(
|
|
uri,
|
|
to,
|
|
ctx=self._ctx,
|
|
include_vectors=include_vectors,
|
|
)
|
|
|
|
async def backup_ovpack(self, to: str, include_vectors: bool = False) -> str:
|
|
"""Back up public scopes as a restore-only .ovpack file."""
|
|
return await self._service.pack.backup_ovpack(
|
|
to,
|
|
ctx=self._ctx,
|
|
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 .ovpack file."""
|
|
return await self._service.pack.import_ovpack(
|
|
file_path,
|
|
parent,
|
|
ctx=self._ctx,
|
|
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 backup .ovpack file."""
|
|
return await self._service.pack.restore_ovpack(
|
|
file_path,
|
|
ctx=self._ctx,
|
|
on_conflict=on_conflict,
|
|
vector_mode=vector_mode,
|
|
)
|
|
|
|
# ============= Git Version Control =============
|
|
|
|
async def git_commit(
|
|
self,
|
|
*,
|
|
message: str,
|
|
paths: Optional[List[str]] = None,
|
|
branch: str = "main",
|
|
author_name: Optional[str] = None,
|
|
author_email: Optional[str] = None,
|
|
) -> Dict[str, Any]:
|
|
"""Create a git snapshot. See VikingFS.commit for semantics."""
|
|
return await self._service.fs.commit(
|
|
message=message,
|
|
paths=paths,
|
|
branch=branch,
|
|
author_name=author_name,
|
|
author_email=author_email,
|
|
ctx=self._ctx,
|
|
)
|
|
|
|
async def git_restore(
|
|
self,
|
|
*,
|
|
project_dir: Optional[str] = None,
|
|
source_commit: str,
|
|
branch: str = "main",
|
|
dry_run: bool = False,
|
|
message: Optional[str] = None,
|
|
author_name: Optional[str] = None,
|
|
author_email: Optional[str] = None,
|
|
) -> Dict[str, Any]:
|
|
"""Restore a subtree, or the full account tree when project_dir is omitted."""
|
|
return await self._service.fs.restore(
|
|
project_dir=project_dir,
|
|
source_commit=source_commit,
|
|
branch=branch,
|
|
dry_run=dry_run,
|
|
message=message,
|
|
author_name=author_name,
|
|
author_email=author_email,
|
|
ctx=self._ctx,
|
|
)
|
|
|
|
async def git_show(
|
|
self,
|
|
target_ref: str,
|
|
*,
|
|
path: Optional[str] = None,
|
|
) -> Any:
|
|
"""Read a commit's metadata or a single blob."""
|
|
return await self._service.fs.show(target_ref, path=path, ctx=self._ctx)
|
|
|
|
async def git_log(
|
|
self,
|
|
*,
|
|
branch: str = "main",
|
|
limit: int = 20,
|
|
paths: Optional[List[str]] = None,
|
|
) -> List[Dict[str, Any]]:
|
|
"""Walk back along parents[0] up to limit commits."""
|
|
return await self._service.fs.log(branch=branch, limit=limit, paths=paths, ctx=self._ctx)
|
|
|
|
async def git_diff(
|
|
self,
|
|
path: str,
|
|
*,
|
|
to_ref: str,
|
|
from_ref: Optional[str] = None,
|
|
) -> Dict[str, Any]:
|
|
"""Compare one file between two snapshot refs."""
|
|
return await self._service.fs.diff(
|
|
path=path,
|
|
from_ref=from_ref,
|
|
to_ref=to_ref,
|
|
ctx=self._ctx,
|
|
)
|
|
|
|
async def git_get_ignore(self) -> str:
|
|
"""Return the account .ovgitignore content (empty string if absent)."""
|
|
return await self._service.fs.get_gitignore(ctx=self._ctx)
|
|
|
|
async def git_set_ignore(self, *, content: str) -> None:
|
|
"""Write the account .ovgitignore control file."""
|
|
await self._service.fs.set_gitignore(content=content, ctx=self._ctx)
|
|
|
|
async def git_delete_ignore(self) -> None:
|
|
"""Delete the account .ovgitignore control file (missing is success)."""
|
|
await self._service.fs.delete_gitignore(ctx=self._ctx)
|
|
|
|
# ============= Debug =============
|
|
|
|
async def check_consistency(self, uri: str) -> Dict[str, Any]:
|
|
"""Check filesystem/vector-index consistency for a URI subtree."""
|
|
return await self._service.check_consistency(
|
|
uri=uri,
|
|
ctx=self._ctx,
|
|
)
|
|
|
|
async def health(self) -> bool:
|
|
"""Check service health."""
|
|
return True # Local service is always healthy if initialized
|
|
|
|
def session(self, session_id: Optional[str] = None, must_exist: bool = False) -> Any:
|
|
"""Create a new session or load an existing one.
|
|
|
|
Args:
|
|
session_id: Session ID, creates a new session if None.
|
|
must_exist: Whether to raise an error if the session does not exist. Default False.
|
|
Returns:
|
|
Session object if exists, None otherwise.
|
|
"""
|
|
|
|
if session_id:
|
|
try:
|
|
return run_async(
|
|
self._service.sessions.get(session_id, self._ctx, auto_create=False)
|
|
)
|
|
except NotFoundError:
|
|
if must_exist:
|
|
raise NotFoundError(session_id, "session")
|
|
|
|
session = self._service.sessions.session(self._ctx, session_id)
|
|
run_async(session.ensure_exists())
|
|
return session
|
|
|
|
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
|
|
"""
|
|
try:
|
|
await self._service.sessions.get(session_id, self._ctx, auto_create=False)
|
|
return True
|
|
except NotFoundError:
|
|
return False
|
|
|
|
def get_status(self) -> Any:
|
|
"""Get system status.
|
|
|
|
Returns:
|
|
SystemStatus containing health status of all components.
|
|
"""
|
|
return self._service.debug.observer.system()
|
|
|
|
def is_healthy(self) -> bool:
|
|
"""Quick health check (synchronous).
|
|
|
|
Returns:
|
|
True if all components are healthy, False otherwise.
|
|
"""
|
|
return self._service.debug.observer.is_healthy()
|
|
|
|
@property
|
|
def observer(self) -> Any:
|
|
"""Get observer service for component status."""
|
|
return self._service.debug.observer
|