Files
OpenViking/openviking/client/session.py
T
DuTao 0ab85f450a feat(session): add turn-aware retention and reliable archive recovery (#3380)
* 优化OpenViking的 session compact逻辑,active message 改为turn,压缩 assistant,保留完整user。
详见RFC:https://github.com/volcengine/OpenViking/discussions/3330

* Vikingbot 使用 ov turn session

* fix pr comment

* 更新文档

* fix pr issue
2026-07-24 14:26:46 +08:00

184 lines
6.2 KiB
Python

# Copyright (c) 2026 Beijing Volcano Engine Technology Co., Ltd.
# SPDX-License-Identifier: AGPL-3.0
"""Lightweight Session class for OpenViking client.
Session delegates all operations to the underlying Client (LocalClient or AsyncHTTPClient).
"""
from dataclasses import asdict
from typing import TYPE_CHECKING, Any, Dict, List, Optional
from openviking.message.part import Part
from openviking.telemetry import TelemetryRequest
from openviking_cli.session.user_id import UserIdentifier
if TYPE_CHECKING:
from openviking_cli.client.base import BaseClient
class Session:
"""Lightweight Session wrapper that delegates operations to Client.
This class provides a convenient OOP interface for session operations.
All actual work is delegated to the underlying client.
"""
def __init__(self, client: "BaseClient", session_id: str, user: UserIdentifier):
"""Initialize Session.
Args:
client: The underlying client (LocalClient or AsyncHTTPClient)
session_id: Session ID
user: User name
"""
self._client = client
self.session_id = session_id
self.user = user
async def add_message(
self,
role: str,
content: Optional[str] = None,
parts: Optional[List[Part]] = None,
created_at: Optional[str] = None,
peer_id: Optional[str] = None,
turn_id: Optional[str] = None,
message_kind: Optional[str] = None,
source_message_ids: Optional[List[str]] = None,
) -> Dict[str, Any]:
"""Add a message to the session.
Args:
role: Message role (e.g., "user", "assistant")
content: Text content (simple mode)
parts: Parts list (TextPart, ContextPart, ImagePart, ToolPart)
created_at: Message creation time (ISO format string). If not provided, current time is used.
peer_id: Optional stable interaction peer identity.
If both content and parts are provided, parts takes precedence.
Returns:
Result dict with session_id and message_count
"""
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
}
if parts is not None:
parts_dicts = [asdict(p) for p in parts]
return await self._client.add_message(
self.session_id,
role,
parts=parts_dicts,
created_at=created_at,
peer_id=peer_id,
**semantic_kwargs,
)
return await self._client.add_message(
self.session_id,
role,
content=content,
created_at=created_at,
peer_id=peer_id,
**semantic_kwargs,
)
async def batch_add_messages(
self,
messages: List[Dict[str, Any]],
) -> Dict[str, Any]:
"""Add multiple messages to the session in a single request.
Args:
messages: List of dicts, each with "role" and optionally "content",
"parts", "created_at", "peer_id".
Returns:
Result dict with session_id, message_count, and added count.
"""
return await self._client.batch_add_messages(
self.session_id,
messages=messages,
)
async def commit(
self,
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 the session (archive messages and extract memories).
Returns:
Commit result
"""
kwargs: Dict[str, Any] = {
"telemetry": telemetry,
"keep_recent_count": keep_recent_count,
}
optional = {
"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,
}
kwargs.update({key: value for key, value in optional.items() if value is not None})
return await self._client.commit_session(self.session_id, **kwargs)
async def commit_async(
self,
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 the session asynchronously (archive messages and extract memories).
Used in viking bot for committing.
Returns:
Commit result
"""
return await self.commit(
telemetry=telemetry,
keep_recent_count=keep_recent_count,
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,
)
async def delete(self) -> None:
"""Delete the session."""
await self._client.delete_session(self.session_id)
async def load(self) -> Dict[str, Any]:
"""Load session data.
Returns:
Session details
"""
return await self._client.get_session(self.session_id)
async def get_session_context(self, token_budget: int = 128_000) -> Dict[str, Any]:
"""Get assembled session context."""
return await self._client.get_session_context(self.session_id, token_budget=token_budget)
async def get_archive(self, archive_id: str) -> Dict[str, Any]:
"""Get one completed archive for the session."""
return await self._client.get_session_archive(self.session_id, archive_id)
def __repr__(self) -> str:
return f"Session(id={self.session_id}, user={self.user.__str__()})"