Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
53 changes: 53 additions & 0 deletions packages/llmpane-py/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -231,6 +231,59 @@ from llmpane.agent import ToolUseMetadata, ToolStatus
# {"metadata": {"tools": [{"name": "get_weather", "status": "completed", "result": "Sunny..."}]}}
```

### Conversation Identity

Conversation IDs are server-owned. Pass an existing ID to continue a
conversation; an unknown or absent ID starts a new one under a
server-generated ID, so a client can never choose the key its data is stored
under. Read the real ID back from the terminal chunk:

```python
async for chunk in session.run(request.conversation_id, request.message):
if chunk.done:
conversation_id = chunk.conversation_id # authoritative
```

A `ConversationStore` is persistence, not authorization — it does not check
*who* may read a conversation. Pass a store already scoped to the current
user, or check ownership before calling `run()`.

### Errors

Stream errors are classified into a stable `ErrorCode`, a safe display
message, and a retryability hint. Raw exception text is **not** sent to the
client: provider exceptions routinely embed prompt content, tool output,
response bodies, or credentials.

The original exception is logged with its traceback through the standard
logging module, so you keep full diagnostics:

```python
import logging

logging.getLogger("llmpane.errors").setLevel(logging.ERROR)
```

For local debugging you can opt back into raw text. Do not enable this in
production — it returns provider exception messages to every client that
triggers an error:

```python
from llmpane.errors import set_expose_raw_errors

set_expose_raw_errors(True) # or LLMPANE_EXPOSE_RAW_ERRORS=1
```

### Terminal Chunk Guarantees

The final chunk is emitted only after the assistant message has been
persisted, and its `message_id` is the ID actually stored — so a client can
rely on that ID existing server-side. If persistence fails, you receive a
classified error instead of a success terminal.

If the client disconnects mid-stream, the user message is already durable and
partial assistant text is discarded rather than silently stored.

## Custom Metadata

Use generics to add type-safe custom metadata:
Expand Down
76 changes: 56 additions & 20 deletions packages/llmpane-py/llmpane/agent/session.py
Original file line number Diff line number Diff line change
Expand Up @@ -3,11 +3,13 @@
from __future__ import annotations

import base64
import uuid
from collections.abc import AsyncGenerator
from typing import TYPE_CHECKING, Any

from llmpane.agent.models import ToolUseMetadata
from llmpane.agent.pydantic_ai import stream_agent
from llmpane.errors import create_error_chunk
from llmpane.limits import ContentLimits, validate_message_content
from llmpane.models import (
ChatMessage,
Expand Down Expand Up @@ -86,8 +88,19 @@ async def run(
) -> AsyncGenerator[StreamChunk[ToolUseMetadata], None]:
"""Run a chat turn with automatic persistence.

Conversation identity is server-owned. An unknown ``conversation_id``
starts a fresh conversation under a server-generated ID rather than
being created under the caller's ID, so a client cannot choose the key
its data is stored under. Read the real ID back from the terminal
chunk's ``conversation_id``.

Note that a ``ConversationStore`` is persistence, not authorization. It
does not check who may read a conversation; pass a store (or a scoped
wrapper) that is already limited to the current user.

Args:
conversation_id: Optional conversation ID (created if None/not found)
conversation_id: Existing conversation to continue. Unknown or
absent values start a new server-generated conversation.
message: User message content (string or list of content parts for multimodal)

Yields:
Expand All @@ -96,12 +109,13 @@ async def run(
if self.content_limits is not None:
validate_message_content(message, self.content_limits)

# Get or create conversation
# Resume only a conversation that already exists. An unknown ID is
# never persisted under the caller-supplied value.
conv = None
if conversation_id:
conv = await self.store.get_conversation(conversation_id)
if not conv:
conv = await self.store.create_conversation(conversation_id)
if conv is None:
conv = await self.store.create_conversation()

# Convert conversation history to Pydantic AI format BEFORE adding new message
history = self._to_pydantic_history(conv.messages)
Expand All @@ -126,30 +140,52 @@ async def run(
limits=self.run_limits,
)
)
# The terminal chunk is held back until persistence has succeeded, so a
# storage failure can never be reported to the client as success. If
# the consumer stops iterating first, the terminal is simply never
# emitted and no partial assistant message is written.
terminal: StreamChunk[ToolUseMetadata] | None = None
async for chunk in stream:
accumulated += chunk.delta
# Add conversation ID to the final chunk
if chunk.done:
yield StreamChunk(
delta=chunk.delta,
metadata=chunk.metadata,
done=True,
error=chunk.error,
error_info=chunk.error_info,
message_id=chunk.message_id,
conversation_id=conv.id,
usage=chunk.usage,
)
else:
yield chunk

# Persist assistant message (after streaming completes)
terminal = chunk
break
yield chunk

if terminal is None:
# The agent stream ended without a terminal event.
yield self._terminal_error(
RuntimeError("Agent stream ended without a final chunk"), conv.id
)
return

if terminal.error or terminal.error_info:
# Already a classified failure; pass it through with identity
# attached and persist nothing.
yield terminal.model_copy(update={"conversation_id": conv.id})
return

# One ID, assigned before persistence, reported after it.
message_id = terminal.message_id or f"msg_{uuid.uuid4().hex[:12]}"
if accumulated:
assistant_msg: ChatMessage[Any] = ChatMessage(
id=message_id,
role=MessageRole.ASSISTANT,
content=accumulated,
)
await self.store.add_message(conv.id, assistant_msg)
try:
await self.store.add_message(conv.id, assistant_msg)
except Exception as exc:
yield self._terminal_error(exc, conv.id)
return

yield terminal.model_copy(update={"message_id": message_id, "conversation_id": conv.id})

@staticmethod
def _terminal_error(exc: Exception, conversation_id: str) -> StreamChunk[Any]:
"""Build a classified terminal error chunk carrying conversation identity."""
chunk = create_error_chunk(exc)
return chunk.model_copy(update={"conversation_id": conversation_id})

def _to_pydantic_prompt(self, content: MessageContent) -> str | list[Any]:
"""Convert llmpane MessageContent to Pydantic AI prompt format.
Expand Down
Loading