# File generated from our OpenAPI spec by Stainless. See CONTRIBUTING.md for details. from __future__ import annotations from typing import Any, Optional, cast from typing_extensions import Literal import httpx from ..._types import Body, Omit, Query, Headers, NotGiven, omit, not_given from ..._utils import path_template, maybe_transform, async_maybe_transform from ..._compat import cached_property from ..._resource import SyncAPIResource, AsyncAPIResource from ..._response import ( to_raw_response_wrapper, to_streamed_response_wrapper, async_to_raw_response_wrapper, async_to_streamed_response_wrapper, ) from ..._streaming import Stream, AsyncStream from ...pagination import SyncArrayPage, AsyncArrayPage from ...types.runs import message_list_params, message_stream_params from ..._base_client import AsyncPaginator, make_request_options from ...types.agents.message import Message from ...types.agents.letta_streaming_response import LettaStreamingResponse __all__ = ["MessagesResource", "AsyncMessagesResource"] class MessagesResource(SyncAPIResource): @cached_property def with_raw_response(self) -> MessagesResourceWithRawResponse: """ This property can be used as a prefix for any HTTP method call to return the raw response object instead of the parsed content. For more information, see https://www.github.com/letta-ai/letta-python#accessing-raw-response-data-eg-headers """ return MessagesResourceWithRawResponse(self) @cached_property def with_streaming_response(self) -> MessagesResourceWithStreamingResponse: """ An alternative to `.with_raw_response` that doesn't eagerly read the response body. For more information, see https://www.github.com/letta-ai/letta-python#with_streaming_response """ return MessagesResourceWithStreamingResponse(self) def list( self, run_id: str, *, after: Optional[str] | Omit = omit, before: Optional[str] | Omit = omit, limit: Optional[int] | Omit = omit, order: Literal["asc", "desc"] | Omit = omit, order_by: Literal["created_at"] | Omit = omit, # Use the following arguments if you need to pass additional parameters to the API that aren't available via kwargs. # The extra values given here take precedence over values defined on the client or passed to this method. extra_headers: Headers | None = None, extra_query: Query | None = None, extra_body: Body | None = None, timeout: float | httpx.Timeout | None | NotGiven = not_given, ) -> SyncArrayPage[Message]: """ Get response messages associated with a run. Args: after: Cursor for pagination (message ID). Returns results relative to this ID in the specified sort order. Expected format: 'message-' before: Cursor for pagination (message ID). Returns results relative to this ID in the specified sort order. Expected format: 'message-' limit: Maximum number of messages to return order: Sort order for messages by creation time. 'asc' for oldest first, 'desc' for newest first order_by: Field to sort by extra_headers: Send extra headers extra_query: Add additional query parameters to the request extra_body: Add additional JSON properties to the request timeout: Override the client-level default timeout for this request, in seconds """ if not run_id: raise ValueError(f"Expected a non-empty value for `run_id` but received {run_id!r}") return self._get_api_list( path_template("/v1/runs/{run_id}/messages", run_id=run_id), page=SyncArrayPage[Message], options=make_request_options( extra_headers=extra_headers, extra_query=extra_query, extra_body=extra_body, timeout=timeout, query=maybe_transform( { "after": after, "before": before, "limit": limit, "order": order, "order_by": order_by, }, message_list_params.MessageListParams, ), ), model=cast(Any, Message), # Union types cannot be passed in as arguments in the type system ) def stream( self, path_run_id: str, *, agent_id: Optional[str] | Omit = omit, batch_size: Optional[int] | Omit = omit, include_pings: Optional[bool] | Omit = omit, otid: Optional[str] | Omit = omit, poll_interval: Optional[float] | Omit = omit, body_run_id: Optional[str] | Omit = omit, starting_after: int | Omit = omit, # Use the following arguments if you need to pass additional parameters to the API that aren't available via kwargs. # The extra values given here take precedence over values defined on the client or passed to this method. extra_headers: Headers | None = None, extra_query: Query | None = None, extra_body: Body | None = None, timeout: float | httpx.Timeout | None | NotGiven = not_given, ) -> Stream[LettaStreamingResponse]: """ Retrieve Stream For Run Args: agent_id: Agent ID for agent-direct mode with 'default' conversation. Use with conversation_id='default' in the URL path. batch_size: Number of entries to read per batch. include_pings: Whether to include periodic keepalive ping messages in the stream to prevent connection timeouts. otid: Offline threading ID to look up the run_id. Bypasses active run lookup if run_id not provided. poll_interval: Seconds to wait between polls when no new data. body_run_id: Run ID to stream directly, bypassing run lookup. Use for recovery from duplicate requests. starting_after: Sequence id to use as a cursor for pagination. Response will start streaming after this chunk sequence id extra_headers: Send extra headers extra_query: Add additional query parameters to the request extra_body: Add additional JSON properties to the request timeout: Override the client-level default timeout for this request, in seconds """ if not path_run_id: raise ValueError(f"Expected a non-empty value for `path_run_id` but received {path_run_id!r}") return self._post( path_template("/v1/runs/{path_run_id}/stream", path_run_id=path_run_id), body=maybe_transform( { "agent_id": agent_id, "batch_size": batch_size, "include_pings": include_pings, "otid": otid, "poll_interval": poll_interval, "body_run_id": body_run_id, "starting_after": starting_after, }, message_stream_params.MessageStreamParams, ), options=make_request_options( extra_headers=extra_headers, extra_query=extra_query, extra_body=extra_body, timeout=timeout ), cast_to=object, stream=True, stream_cls=Stream[LettaStreamingResponse], ) class AsyncMessagesResource(AsyncAPIResource): @cached_property def with_raw_response(self) -> AsyncMessagesResourceWithRawResponse: """ This property can be used as a prefix for any HTTP method call to return the raw response object instead of the parsed content. For more information, see https://www.github.com/letta-ai/letta-python#accessing-raw-response-data-eg-headers """ return AsyncMessagesResourceWithRawResponse(self) @cached_property def with_streaming_response(self) -> AsyncMessagesResourceWithStreamingResponse: """ An alternative to `.with_raw_response` that doesn't eagerly read the response body. For more information, see https://www.github.com/letta-ai/letta-python#with_streaming_response """ return AsyncMessagesResourceWithStreamingResponse(self) def list( self, run_id: str, *, after: Optional[str] | Omit = omit, before: Optional[str] | Omit = omit, limit: Optional[int] | Omit = omit, order: Literal["asc", "desc"] | Omit = omit, order_by: Literal["created_at"] | Omit = omit, # Use the following arguments if you need to pass additional parameters to the API that aren't available via kwargs. # The extra values given here take precedence over values defined on the client or passed to this method. extra_headers: Headers | None = None, extra_query: Query | None = None, extra_body: Body | None = None, timeout: float | httpx.Timeout | None | NotGiven = not_given, ) -> AsyncPaginator[Message, AsyncArrayPage[Message]]: """ Get response messages associated with a run. Args: after: Cursor for pagination (message ID). Returns results relative to this ID in the specified sort order. Expected format: 'message-' before: Cursor for pagination (message ID). Returns results relative to this ID in the specified sort order. Expected format: 'message-' limit: Maximum number of messages to return order: Sort order for messages by creation time. 'asc' for oldest first, 'desc' for newest first order_by: Field to sort by extra_headers: Send extra headers extra_query: Add additional query parameters to the request extra_body: Add additional JSON properties to the request timeout: Override the client-level default timeout for this request, in seconds """ if not run_id: raise ValueError(f"Expected a non-empty value for `run_id` but received {run_id!r}") return self._get_api_list( path_template("/v1/runs/{run_id}/messages", run_id=run_id), page=AsyncArrayPage[Message], options=make_request_options( extra_headers=extra_headers, extra_query=extra_query, extra_body=extra_body, timeout=timeout, query=maybe_transform( { "after": after, "before": before, "limit": limit, "order": order, "order_by": order_by, }, message_list_params.MessageListParams, ), ), model=cast(Any, Message), # Union types cannot be passed in as arguments in the type system ) async def stream( self, path_run_id: str, *, agent_id: Optional[str] | Omit = omit, batch_size: Optional[int] | Omit = omit, include_pings: Optional[bool] | Omit = omit, otid: Optional[str] | Omit = omit, poll_interval: Optional[float] | Omit = omit, body_run_id: Optional[str] | Omit = omit, starting_after: int | Omit = omit, # Use the following arguments if you need to pass additional parameters to the API that aren't available via kwargs. # The extra values given here take precedence over values defined on the client or passed to this method. extra_headers: Headers | None = None, extra_query: Query | None = None, extra_body: Body | None = None, timeout: float | httpx.Timeout | None | NotGiven = not_given, ) -> AsyncStream[LettaStreamingResponse]: """ Retrieve Stream For Run Args: agent_id: Agent ID for agent-direct mode with 'default' conversation. Use with conversation_id='default' in the URL path. batch_size: Number of entries to read per batch. include_pings: Whether to include periodic keepalive ping messages in the stream to prevent connection timeouts. otid: Offline threading ID to look up the run_id. Bypasses active run lookup if run_id not provided. poll_interval: Seconds to wait between polls when no new data. body_run_id: Run ID to stream directly, bypassing run lookup. Use for recovery from duplicate requests. starting_after: Sequence id to use as a cursor for pagination. Response will start streaming after this chunk sequence id extra_headers: Send extra headers extra_query: Add additional query parameters to the request extra_body: Add additional JSON properties to the request timeout: Override the client-level default timeout for this request, in seconds """ if not path_run_id: raise ValueError(f"Expected a non-empty value for `path_run_id` but received {path_run_id!r}") return await self._post( path_template("/v1/runs/{path_run_id}/stream", path_run_id=path_run_id), body=await async_maybe_transform( { "agent_id": agent_id, "batch_size": batch_size, "include_pings": include_pings, "otid": otid, "poll_interval": poll_interval, "body_run_id": body_run_id, "starting_after": starting_after, }, message_stream_params.MessageStreamParams, ), options=make_request_options( extra_headers=extra_headers, extra_query=extra_query, extra_body=extra_body, timeout=timeout ), cast_to=object, stream=True, stream_cls=AsyncStream[LettaStreamingResponse], ) class MessagesResourceWithRawResponse: def __init__(self, messages: MessagesResource) -> None: self._messages = messages self.list = to_raw_response_wrapper( messages.list, ) self.stream = to_raw_response_wrapper( messages.stream, ) class AsyncMessagesResourceWithRawResponse: def __init__(self, messages: AsyncMessagesResource) -> None: self._messages = messages self.list = async_to_raw_response_wrapper( messages.list, ) self.stream = async_to_raw_response_wrapper( messages.stream, ) class MessagesResourceWithStreamingResponse: def __init__(self, messages: MessagesResource) -> None: self._messages = messages self.list = to_streamed_response_wrapper( messages.list, ) self.stream = to_streamed_response_wrapper( messages.stream, ) class AsyncMessagesResourceWithStreamingResponse: def __init__(self, messages: AsyncMessagesResource) -> None: self._messages = messages self.list = async_to_streamed_response_wrapper( messages.list, ) self.stream = async_to_streamed_response_wrapper( messages.stream, )