|
| 1 | +from __future__ import annotations |
| 2 | +import enum |
| 3 | +import json |
| 4 | +import re |
| 5 | +import sys |
| 6 | +from datetime import datetime, timezone |
| 7 | +from typing import Optional, List, Dict, Any |
| 8 | +from urllib.parse import quote |
| 9 | + |
| 10 | +from ravendb.documents.operations.definitions import MaintenanceOperation |
| 11 | +from ravendb.documents.conventions import DocumentConventions |
| 12 | +from ravendb.http.raven_command import RavenCommand |
| 13 | +from ravendb.http.server_node import ServerNode |
| 14 | +import requests |
| 15 | + |
| 16 | + |
| 17 | +class AiConversationDetailLevel(enum.Enum): |
| 18 | + """Controls the level of detail when reading conversation messages.""" |
| 19 | + |
| 20 | + SIMPLE = "Simple" |
| 21 | + DETAILED = "Detailed" |
| 22 | + FULL = "Full" |
| 23 | + |
| 24 | + def __str__(self) -> str: |
| 25 | + return self.value |
| 26 | + |
| 27 | + |
| 28 | +class AiMessageRole(enum.Enum): |
| 29 | + """The role of a message sender in an AI conversation.""" |
| 30 | + |
| 31 | + SYSTEM = "System" |
| 32 | + USER = "User" |
| 33 | + ASSISTANT = "Assistant" |
| 34 | + SUMMARY = "Summary" |
| 35 | + INTERNAL = "Internal" |
| 36 | + |
| 37 | + def __str__(self) -> str: |
| 38 | + return self.value |
| 39 | + |
| 40 | + |
| 41 | +class AiToolCallResult: |
| 42 | + """Represents a tool call result from an AI conversation message.""" |
| 43 | + |
| 44 | + def __init__( |
| 45 | + self, |
| 46 | + id: Optional[str] = None, |
| 47 | + name: Optional[str] = None, |
| 48 | + arguments: Optional[str] = None, |
| 49 | + result: Optional[str] = None, |
| 50 | + sub_conversation_id: Optional[str] = None, |
| 51 | + ): |
| 52 | + self.id: Optional[str] = id |
| 53 | + self.name: Optional[str] = name |
| 54 | + self.arguments: Optional[str] = arguments |
| 55 | + self.result: Optional[str] = result |
| 56 | + self.sub_conversation_id: Optional[str] = sub_conversation_id |
| 57 | + |
| 58 | + @classmethod |
| 59 | + def from_json(cls, json_dict: Dict[str, Any]) -> AiToolCallResult: |
| 60 | + return cls( |
| 61 | + id=json_dict.get("Id"), |
| 62 | + name=json_dict.get("Name"), |
| 63 | + arguments=json_dict.get("Arguments"), |
| 64 | + result=json_dict.get("Result"), |
| 65 | + sub_conversation_id=json_dict.get("SubConversationId"), |
| 66 | + ) |
| 67 | + |
| 68 | + def to_json(self) -> Dict[str, Any]: |
| 69 | + return { |
| 70 | + "Id": self.id, |
| 71 | + "Name": self.name, |
| 72 | + "Arguments": self.arguments, |
| 73 | + "Result": self.result, |
| 74 | + "SubConversationId": self.sub_conversation_id, |
| 75 | + } |
| 76 | + |
| 77 | + |
| 78 | +class AiConversationMessage: |
| 79 | + """Represents a single message in an AI conversation.""" |
| 80 | + |
| 81 | + def __init__( |
| 82 | + self, |
| 83 | + role: Optional[AiMessageRole] = None, |
| 84 | + content: Optional[str] = None, |
| 85 | + attachments: Optional[List[str]] = None, |
| 86 | + timestamp: Optional[datetime] = None, |
| 87 | + tool_calls: Optional[List[AiToolCallResult]] = None, |
| 88 | + usage: Optional[Any] = None, |
| 89 | + sub_conversation_id: Optional[str] = None, |
| 90 | + ): |
| 91 | + self.role: Optional[AiMessageRole] = role |
| 92 | + self.content: Optional[str] = content |
| 93 | + self.attachments: List[str] = attachments or [] |
| 94 | + self.timestamp: Optional[datetime] = timestamp |
| 95 | + self.tool_calls: List[AiToolCallResult] = tool_calls or [] |
| 96 | + self.usage: Optional[Any] = usage |
| 97 | + self.sub_conversation_id: Optional[str] = sub_conversation_id |
| 98 | + |
| 99 | + @classmethod |
| 100 | + def from_json(cls, json_dict: Dict[str, Any]) -> AiConversationMessage: |
| 101 | + from ravendb.documents.operations.ai.agents.run_conversation_operation import AiUsage |
| 102 | + |
| 103 | + return cls( |
| 104 | + role=AiMessageRole(json_dict["Role"]) if json_dict.get("Role") else None, |
| 105 | + content=json_dict.get("Content"), |
| 106 | + attachments=json_dict.get("Attachments") or [], |
| 107 | + timestamp=_parse_raven_datetime(json_dict.get("Timestamp")), |
| 108 | + tool_calls=( |
| 109 | + [AiToolCallResult.from_json(tc) for tc in json_dict["ToolCalls"]] |
| 110 | + if json_dict.get("ToolCalls") |
| 111 | + else [] |
| 112 | + ), |
| 113 | + usage=AiUsage.from_json(json_dict["Usage"]) if json_dict.get("Usage") else None, |
| 114 | + sub_conversation_id=json_dict.get("SubConversationId"), |
| 115 | + ) |
| 116 | + |
| 117 | + def to_json(self) -> Dict[str, Any]: |
| 118 | + return { |
| 119 | + "Role": self.role.value if self.role else None, |
| 120 | + "Content": self.content, |
| 121 | + "Attachments": self.attachments, |
| 122 | + "Timestamp": _format_raven_datetime(self.timestamp) if self.timestamp else None, |
| 123 | + "ToolCalls": [tc.to_json() for tc in self.tool_calls] if self.tool_calls else None, |
| 124 | + "Usage": self.usage.to_json() if self.usage else None, |
| 125 | + "SubConversationId": self.sub_conversation_id, |
| 126 | + } |
| 127 | + |
| 128 | + |
| 129 | +class AiConversationMessagesResult: |
| 130 | + """The result of fetching conversation messages.""" |
| 131 | + |
| 132 | + def __init__( |
| 133 | + self, |
| 134 | + conversation_id: Optional[str] = None, |
| 135 | + agent: Optional[str] = None, |
| 136 | + parameters: Optional[Dict[str, Any]] = None, |
| 137 | + total_usage: Optional[Any] = None, |
| 138 | + last_message_at: Optional[datetime] = None, |
| 139 | + messages: Optional[List[AiConversationMessage]] = None, |
| 140 | + has_more_messages: bool = False, |
| 141 | + sub_conversation_ids: Optional[List[str]] = None, |
| 142 | + attachments: Optional[List[str]] = None, |
| 143 | + ): |
| 144 | + self.conversation_id: Optional[str] = conversation_id |
| 145 | + self.agent: Optional[str] = agent |
| 146 | + self.parameters: Dict[str, Any] = parameters or {} |
| 147 | + self.total_usage: Optional[Any] = total_usage |
| 148 | + self.last_message_at: Optional[datetime] = last_message_at |
| 149 | + self.messages: List[AiConversationMessage] = messages or [] |
| 150 | + self.has_more_messages: bool = has_more_messages |
| 151 | + self.sub_conversation_ids: List[str] = sub_conversation_ids or [] |
| 152 | + self.attachments: List[str] = attachments or [] |
| 153 | + |
| 154 | + @classmethod |
| 155 | + def from_json(cls, json_dict: Dict[str, Any]) -> AiConversationMessagesResult: |
| 156 | + from ravendb.documents.operations.ai.agents.run_conversation_operation import AiUsage |
| 157 | + |
| 158 | + return cls( |
| 159 | + conversation_id=json_dict.get("ConversationId"), |
| 160 | + agent=json_dict.get("Agent"), |
| 161 | + parameters=json_dict.get("Parameters") or {}, |
| 162 | + total_usage=AiUsage.from_json(json_dict["TotalUsage"]) if json_dict.get("TotalUsage") else None, |
| 163 | + last_message_at=_parse_raven_datetime(json_dict.get("LastMessageAt")), |
| 164 | + messages=( |
| 165 | + [AiConversationMessage.from_json(msg) for msg in json_dict["Messages"]] |
| 166 | + if json_dict.get("Messages") |
| 167 | + else [] |
| 168 | + ), |
| 169 | + has_more_messages=json_dict.get("HasMoreMessages", False), |
| 170 | + sub_conversation_ids=json_dict.get("SubConversationIds") or [], |
| 171 | + attachments=json_dict.get("Attachments") or [], |
| 172 | + ) |
| 173 | + |
| 174 | + def to_json(self) -> Dict[str, Any]: |
| 175 | + return { |
| 176 | + "ConversationId": self.conversation_id, |
| 177 | + "Agent": self.agent, |
| 178 | + "Parameters": self.parameters, |
| 179 | + "TotalUsage": self.total_usage.to_json() if self.total_usage else None, |
| 180 | + "LastMessageAt": _format_raven_datetime(self.last_message_at) if self.last_message_at else None, |
| 181 | + "Messages": [msg.to_json() for msg in self.messages] if self.messages else None, |
| 182 | + "HasMoreMessages": self.has_more_messages, |
| 183 | + "SubConversationIds": self.sub_conversation_ids, |
| 184 | + "Attachments": self.attachments, |
| 185 | + } |
| 186 | + |
| 187 | + |
| 188 | +class GetConversationMessagesOptions: |
| 189 | + """Parameters for reading messages from an AI agent conversation.""" |
| 190 | + |
| 191 | + def __init__( |
| 192 | + self, |
| 193 | + conversation_id: str, |
| 194 | + before: Optional[datetime] = None, |
| 195 | + after: Optional[datetime] = None, |
| 196 | + page_size: int = sys.maxsize, |
| 197 | + detail_level: AiConversationDetailLevel = AiConversationDetailLevel.SIMPLE, |
| 198 | + ): |
| 199 | + if not conversation_id or conversation_id.isspace(): |
| 200 | + raise ValueError("conversation_id cannot be None or empty") |
| 201 | + |
| 202 | + if before is not None and after is not None: |
| 203 | + raise ValueError("before and after cannot both be specified") |
| 204 | + |
| 205 | + if page_size <= 0: |
| 206 | + raise ValueError("page_size must be greater than 0") |
| 207 | + |
| 208 | + self.conversation_id: str = conversation_id |
| 209 | + self.before: Optional[datetime] = before |
| 210 | + self.after: Optional[datetime] = after |
| 211 | + self.page_size: int = page_size |
| 212 | + self.detail_level: AiConversationDetailLevel = detail_level |
| 213 | + |
| 214 | + |
| 215 | +class GetConversationMessagesOperation(MaintenanceOperation[AiConversationMessagesResult]): |
| 216 | + """ |
| 217 | + Reads messages from an AI agent conversation, with optional timestamp-based paging and view filtering. |
| 218 | + """ |
| 219 | + |
| 220 | + def __init__(self, conversation_id_or_options: str | GetConversationMessagesOptions): |
| 221 | + if isinstance(conversation_id_or_options, str): |
| 222 | + self._parameters = GetConversationMessagesOptions(conversation_id=conversation_id_or_options) |
| 223 | + elif isinstance(conversation_id_or_options, GetConversationMessagesOptions): |
| 224 | + self._parameters = conversation_id_or_options |
| 225 | + else: |
| 226 | + raise TypeError( |
| 227 | + "Expected str (conversation_id) or GetConversationMessagesOptions, " |
| 228 | + f"got {type(conversation_id_or_options).__name__}" |
| 229 | + ) |
| 230 | + |
| 231 | + def get_command(self, conventions: DocumentConventions) -> RavenCommand[AiConversationMessagesResult]: |
| 232 | + return GetConversationMessagesCommand(self._parameters) |
| 233 | + |
| 234 | + |
| 235 | +class GetConversationMessagesCommand(RavenCommand[AiConversationMessagesResult]): |
| 236 | + def __init__(self, parameters: GetConversationMessagesOptions): |
| 237 | + super().__init__(AiConversationMessagesResult) |
| 238 | + self._params = parameters |
| 239 | + |
| 240 | + def is_read_request(self) -> bool: |
| 241 | + return True |
| 242 | + |
| 243 | + def create_request(self, node: ServerNode) -> requests.Request: |
| 244 | + url = ( |
| 245 | + f"{node.url}/databases/{node.database}/ai/agent/conversation/messages" |
| 246 | + f"?conversationId={quote(self._params.conversation_id, safe='')}" |
| 247 | + ) |
| 248 | + |
| 249 | + if self._params.before is not None: |
| 250 | + url += f"&before={quote(_format_raven_datetime(self._params.before), safe='')}" |
| 251 | + |
| 252 | + if self._params.after is not None: |
| 253 | + url += f"&after={quote(_format_raven_datetime(self._params.after), safe='')}" |
| 254 | + |
| 255 | + url += f"&pageSize={self._params.page_size}" |
| 256 | + url += f"&detailLevel={self._params.detail_level.value}" |
| 257 | + |
| 258 | + return requests.Request("GET", url) |
| 259 | + |
| 260 | + def set_response(self, response: str, from_cache: bool) -> None: |
| 261 | + if response is None: |
| 262 | + self.result = None |
| 263 | + return |
| 264 | + |
| 265 | + response_json = json.loads(response) |
| 266 | + self.result = AiConversationMessagesResult.from_json(response_json) |
| 267 | + |
| 268 | + |
| 269 | +# ------------------------------------------------------------------ |
| 270 | +# Raven datetime helpers — wire contract uses yyyy-MM-ddTHH:mm:ss.fffffffZ |
| 271 | +# ------------------------------------------------------------------ |
| 272 | + |
| 273 | + |
| 274 | +def _ensure_utc(dt: datetime) -> datetime: |
| 275 | + """Ensure a datetime is in UTC.""" |
| 276 | + if dt.tzinfo is None: |
| 277 | + return dt.replace(tzinfo=timezone.utc) |
| 278 | + return dt.astimezone(timezone.utc) |
| 279 | + |
| 280 | + |
| 281 | +def _format_raven_datetime(dt: datetime) -> str: |
| 282 | + """Format a datetime in Raven's default UTC format: yyyy-MM-ddTHH:mm:ss.fffffffZ.""" |
| 283 | + dt = _ensure_utc(dt) |
| 284 | + return dt.strftime("%Y-%m-%dT%H:%M:%S.%f") + "0Z" |
| 285 | + |
| 286 | + |
| 287 | +def _parse_raven_datetime(value: Optional[str]) -> Optional[datetime]: |
| 288 | + """Parse a Raven datetime string (e.g. 2025-06-15T12:00:00.0000000Z) into a UTC datetime.""" |
| 289 | + if value is None: |
| 290 | + return None |
| 291 | + |
| 292 | + if value.endswith("Z"): |
| 293 | + value = value[:-1] + "+00:00" |
| 294 | + |
| 295 | + # fromisoformat before Python 3.11 accepts only 3 or 6 fractional digits; Raven emits 7 (100ns ticks). |
| 296 | + value = re.sub(r"\.(\d+)", lambda m: "." + (m.group(1) + "000000")[:6], value) |
| 297 | + |
| 298 | + return datetime.fromisoformat(value) |
0 commit comments