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
38 changes: 26 additions & 12 deletions docs/agent-loop.md
Original file line number Diff line number Diff line change
Expand Up @@ -183,6 +183,25 @@ Do not call `asyncio.run(...)` from inside an already-running event loop
runtime). In those environments, keep your Pollux adapter async and `await`
`interact()`, `run()`, or `stream()` directly.

### Closing streams early

Cancellation of a task awaiting `interact()` or a stream iteration propagates
through Pollux and releases the active provider request. If a consumer may stop
streaming before the terminal `done` event, explicitly close the iterator; a
bare `break` does not guarantee synchronous async-generator cleanup:

```python
from contextlib import aclosing

async with aclosing(stream(env, input, config=config)) as events:
async for event in events:
if should_stop(event):
break
```

Closing a `Session.stream()` releases only that request. The session remains
open and reusable until its own context manager exits or `aclose()` is called.

## Variations

These are all small modifications to the same loop structure.
Expand Down Expand Up @@ -234,27 +253,22 @@ out = await interact(
```

If your application owns an OpenAI Chat Completions-style transcript for resume,
compaction, or audit logs, keep that transcript as the durable record and import
it into a fresh Pollux continuation for each turn:
compaction, or audit logs, convert its portable messages into typed history:

```python
from pollux import Continuation, Input
from pollux import Input, Message

continuation = Continuation.from_openai_messages(messages, provider="local")
history = [Message.from_openai(message) for message in messages]
out = await interact(
env,
Input(content="Continue.", continuation=continuation),
Input(content="Continue.", history=history),
config=config,
)
```

`ToolCall.to_openai()`, `Message.to_openai()`, and
`Continuation.to_openai_messages()` provide the reverse mapping for harnesses
that still dispatch OpenAI-shaped tool calls. For supported text and tool
message shapes, importing with `Continuation.from_openai_messages(...)` and
exporting with `to_openai_messages()` preserves `role`, `content`,
`tool_calls`, and `tool_call_id`. Display reasoning is output data and is not
replayed as transcript history.
`ToolCall.to_openai()` and `Message.to_openai()` provide the reverse mapping for
harnesses that dispatch OpenAI-shaped tool calls. System instructions belong on
`Environment`, and display reasoning is not replayed as transcript history.

### Guiding tool use with system instructions

Expand Down
64 changes: 38 additions & 26 deletions docs/conversations-and-agents.md
Original file line number Diff line number Diff line change
Expand Up @@ -32,7 +32,11 @@ next turn.

## Continuing a Conversation with `Continuation`

Pass a prior result's `continuation` back into the next `Input(continuation=...)` to automatically resume a conversation. Pollux unpacks the initial prompt, the assistant's previous response, and any tool calls directly into the context payload.
Pass a prior result's `continuation` back into the next
`Input(continuation=...)` to resume the provider-correct conversation. Treat the
value as opaque: serialize it with `to_jsonable()`, restore it with
`Continuation.from_jsonable()`, and otherwise pass it back unchanged. Do not
inspect, edit, merge, or summarize its serialized provider state.

To get a `continuation` for subsequent turns in plain conversational calls, the first turn must opt into conversation tracking by passing `history=[]` (or an empty list/tuple). Without it, Pollux treats the call as stateless and does not build continuation state.

Expand Down Expand Up @@ -74,8 +78,7 @@ If you need to inject mid-conversation context, groom old context out to
save tokens, or resume a chat from a database, a prior `continuation` alone
is not enough.

Instead, pass an explicit `history` list of dictionaries containing `role`
and `content`:
Instead, pass explicit typed `Message` history:

```python
import asyncio
Expand Down Expand Up @@ -104,7 +107,9 @@ asyncio.run(manual_history_injection())
```

Pollux treats the `history` block chronologically *before* the prompt you
provide to `interact()`.
provide to `interact()`. This deliberately gives up response IDs and opaque
provider replay state. A successful interaction creates a fresh continuation
for the active provider.

## Persisting Agent Transcripts

Expand All @@ -123,42 +128,43 @@ product transcript:

If your application is the source of truth for history because it supports
resume, compaction, truncation, or audit logs, keep that transcript in your own
store and rebuild a `Continuation` for each Pollux turn. In this pattern,
`Continuation` is the provider replay object for the next call, not the durable
application transcript.

`Continuation.from_openai_messages(...)` is a compatibility bridge for text
Chat Completions-style transcripts and tool turns. It extracts text from
OpenAI text parts, preserves tool calls, and deliberately does not turn
provider-shaped media attachments into Pollux `Source` objects. That keeps
media explicit at the Pollux boundary:
store and rebuild typed `Message` history for each Pollux turn:

```python
from pollux import Continuation, Environment, Input, Source, interact
from pollux import Environment, Input, Message, Source, interact

continuation = Continuation.from_openai_messages(text_messages, provider="local")
history = [Message.from_openai(message) for message in text_messages]
env = Environment(sources=[Source.from_file("current-report.pdf")])

result = await interact(
env,
Input(content="Continue the analysis.", continuation=continuation),
Input(content="Continue the analysis.", history=history),
config=config,
)
```

For supported text and tool message shapes,
`Continuation.from_openai_messages(messages).to_openai_messages()` is a
lossless round trip for `role`, `content`, `tool_calls`, and `tool_call_id`.
Display reasoning is not part of that replay contract: `Output.reasoning` is
diagnostic/output data and should not be copied into future transcript messages.
Provider-specific opaque reasoning blocks are replayed through
`provider_state` only when a provider requires them.
`Message.from_openai()` and `Message.to_openai()` bridge individual portable
text/tool messages. Move OpenAI `system` messages to `Environment.instructions`.
Display reasoning is diagnostic output and should not be copied into transcript
messages.

If you need to resume an older session that included files or images, load
those application records and rebuild `Source.from_file(...)`,
`Source.from_uri(...)`, or another source constructor explicitly. Pollux's
replay messages are for conversation turns, not hidden media transport.

### Portable History Shapes

Portable history consists of non-empty user text, assistant text and/or
normalized `ToolCall` values, and tool messages with the matching
`tool_call_id`. Keep parallel tool results in call order. System messages,
media, reasoning, response IDs, and provider-native state are not portable.

Anthropic extended-thinking tool turns are a special boundary: signed thinking
blocks live only in the untouched continuation. Return pending tool results with
that continuation, then compact at a completed interaction boundary. Groomed
history cannot recreate the signatures Anthropic requires.

## Handling Tool Messages in History

If your conversation includes tool execution (the model asked for data, you
Expand Down Expand Up @@ -256,14 +262,20 @@ next_result = await interact(
- **Reasoning is output data unless a provider requires opaque replay.**
Pollux surfaces reasoning text on `Output.reasoning` for display and
debugging. Provider-specific signed thinking or reasoning blocks are kept
inside continuation `provider_state` only when the provider requires them
for valid follow-up turns; do not copy display reasoning into future user or
assistant messages yourself.
inside the opaque continuation only when the provider requires them for valid
follow-up turns; do not copy display reasoning into future messages.
- **Provider differences exist.** Gemini, OpenAI, and Anthropic support tool
calling and tool messages in history. OpenRouter supports them on models
that advertise tool support. See
[Provider Capabilities](reference/provider-capabilities.md) for details.

## Durable Environment Identity

Use `environment.fingerprint(provider=config.provider)` to bind saved work to
the model-facing instructions, sources, tools, and provider. Compose
`config.model` plus application policy/schema identifiers separately. Cache
preferences and environment metadata intentionally do not affect this identity.

---

Now that you understand the conversation mechanics, see
Expand Down
16 changes: 11 additions & 5 deletions docs/migrating-to-v2.md
Original file line number Diff line number Diff line change
Expand Up @@ -104,11 +104,17 @@ produced them, and `from_jsonable()` rejects an incompatible version (and, when
pass `expected_provider=`, a mismatched provider) with an actionable error instead
of misreading it.

A continuation is bound to its provider: its `provider_state` (response ids,
provider-specific replay blocks) is not portable, so reusing one under a different
provider is rejected before dispatch. Across the 1.x → 2.0 boundary, plan to re-run
work rather than reusing old serialized blobs; persist enough application state to
rebuild the request when that is the right recovery path.
A continuation is opaque and bound to its provider. Applications serialize,
restore, and pass it back unchanged; provider response IDs and replay blocks are
not editable or portable. The final v2 RC uses continuation schema version 2 and
rejects schema-v1 RC artifacts. `Continuation.from_openai_messages()` and
`to_openai_messages()` were removed; use typed `Message` history after grooming
or importing an application transcript.

For durable run identity, use
`environment.fingerprint(provider=config.provider)` and compose `config.model`
plus application policy/schema versions separately. `EnvironmentSnapshot` is
now internal and is no longer a supported identity route.

## What To Do In 1.x

Expand Down
2 changes: 2 additions & 0 deletions docs/reference/api.md
Original file line number Diff line number Diff line change
Expand Up @@ -88,6 +88,8 @@ and `ResultEnvelope` types are no longer part of the public API.

::: pollux.Continuation

::: pollux.Message

::: pollux.ToolDeclaration

::: pollux.ToolCall
Expand Down
3 changes: 3 additions & 0 deletions docs/reference/provider-capabilities.md
Original file line number Diff line number Diff line change
Expand Up @@ -144,6 +144,9 @@ jobs, see [Building With Deferred Delivery](../building-with-deferred-delivery.m
holds for `stream()` too: signed thinking blocks are reassembled from the
stream, so a streamed extended-thinking + tool turn continues identically to
the non-streaming path.
- Do not compact between an extended-thinking tool call and its tool results.
The signed blocks exist only in the opaque continuation; complete the tool
exchange first, then switch to application-authored `Message` history.
- `max_tokens`: limits the output length. Default is `16384` for Anthropic,
which leaves room for thinking output at all effort levels. Other providers
currently ignore this option.
Expand Down
33 changes: 22 additions & 11 deletions src/pollux/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,7 @@
import logging
from typing import TYPE_CHECKING, Any, cast

from pollux._lifecycle import close_async_iterator
from pollux.config import (
_API_KEY_ENV_VARS,
_LOCAL_BASE_URL_ENV_VAR,
Expand Down Expand Up @@ -53,10 +54,10 @@
CacheSetting,
Continuation,
Environment,
EnvironmentSnapshot,
Event,
Input,
Message,
MessageRole,
Output,
OutputCollection,
OutputRequirements,
Expand All @@ -66,6 +67,7 @@
ToolDeclaration,
ToolResult,
)
from pollux.interaction.environment import EnvironmentSnapshot as _EnvironmentSnapshot
from pollux.interaction.execute import (
execute_interaction,
execute_interactions,
Expand All @@ -81,7 +83,7 @@
from pollux.source import Source

if TYPE_CHECKING:
from collections.abc import AsyncIterator, Callable, Sequence
from collections.abc import AsyncGenerator, Callable, Sequence

from pollux.interaction.schema import ResponseSchemaInput
from pollux.providers.base import Provider
Expand Down Expand Up @@ -287,7 +289,7 @@ async def stream(
reasoning_budget_tokens: int | None = None,
tool_choice: ToolChoice | None = None,
provider_options: dict[str, dict[str, Any]] | None = None,
) -> AsyncIterator[Event]:
) -> AsyncGenerator[Event, None]:
"""Stream one explicit v2 interaction as :class:`Event` objects.

The streaming sibling of :func:`interact`: same environment/input/config and
Expand Down Expand Up @@ -326,7 +328,7 @@ async def stream(
result = event.output
"""
async with Session(config) as session:
async for event in session.stream(
events = session.stream(
environment,
input,
output=output,
Expand All @@ -338,8 +340,12 @@ async def stream(
reasoning_budget_tokens=reasoning_budget_tokens,
tool_choice=tool_choice,
provider_options=provider_options,
):
yield event
)
try:
async for event in events:
yield event
finally:
await close_async_iterator(events)


class Session:
Expand Down Expand Up @@ -421,7 +427,7 @@ async def stream(
reasoning_budget_tokens: int | None = None,
tool_choice: ToolChoice | None = None,
provider_options: dict[str, dict[str, Any]] | None = None,
) -> AsyncIterator[Event]:
) -> AsyncGenerator[Event, None]:
"""Stream one interaction using the session's provider instance."""
self._ensure_open()
requirements = _build_requirements(
Expand All @@ -435,10 +441,14 @@ async def stream(
tool_choice=tool_choice,
provider_options=provider_options,
)
async for event in stream_interaction(
events = stream_interaction(
environment, input, requirements, self.config, self._provider
):
yield event
)
try:
async for event in events:
yield event
finally:
await close_async_iterator(events)

async def run_many(
self,
Expand Down Expand Up @@ -733,7 +743,7 @@ async def prepare_environment(
if isinstance(cache, CachePolicy):
provider = _get_provider(config)
try:
snapshot = EnvironmentSnapshot.from_environment(
snapshot = _EnvironmentSnapshot.from_environment(
environment, provider=config.provider
)
await resolve_persistent_cache(snapshot, config, provider)
Expand Down Expand Up @@ -965,6 +975,7 @@ def _resolve_deferred_provider(handle: DeferredHandle) -> Provider:
"Input",
"InternalError",
"Message",
"MessageRole",
"Output",
"OutputCollection",
"OutputRequirements",
Expand Down
43 changes: 43 additions & 0 deletions src/pollux/_lifecycle.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,43 @@
"""Shared asynchronous resource-lifecycle helpers."""

from __future__ import annotations

import asyncio
import inspect
import logging
from typing import Any
from weakref import WeakSet

logger = logging.getLogger(__name__)
_closed_iterators: WeakSet[Any] = WeakSet()


async def close_async_iterator(iterator: Any) -> None:
"""Close an async iterator without masking an interaction's primary error."""
close = getattr(iterator, "aclose", None)
if not callable(close):
close = getattr(iterator, "close", None)
if not callable(close):
return
try:
if iterator in _closed_iterators:
return
_closed_iterators.add(iterator)
except TypeError:
# Some third-party stream wrappers cannot be weak-referenced. Prefer a
# private marker when they allow attributes; otherwise rely on the SDK's
# idempotent close contract.
try:
if getattr(iterator, "_pollux_closed", False):
return
iterator._pollux_closed = True
except (AttributeError, TypeError):
pass
try:
result = close()
if inspect.isawaitable(result):
await result
except asyncio.CancelledError:
raise
except Exception as exc:
logger.warning("Async iterator cleanup failed: %s", exc)
Loading
Loading