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
2 changes: 1 addition & 1 deletion AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -131,7 +131,7 @@ Important Pi RPC assumptions:
- Start command is `pi --mode rpc`.
- Existing sessions can resume with `--session <session-file>`.
- JSONL records are newline-delimited.
- Prompt completion is detected by consuming events until `agent_end`.
- Prompt completion requires `agent_settled`; `agent_end` only ends one low-level run and Pi may still compact/retry.

If Pi RPC protocol changes, update:

Expand Down
2 changes: 1 addition & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -183,7 +183,7 @@ uv run pi-gateway run
- `/steer <text>` steer current/next turn
- `/pi <text>` send raw text to Pi, including Pi slash commands

Normal Telegram messages are sent to Pi as prompts.
Normal Telegram messages are sent to Pi as prompts. Pi handles automatic context compaction (when enabled in Pi settings); the gateway waits for Pi's `agent_settled` event before replying, including any overflow recovery, retries, or queued work after an `agent_end`. This requires a Pi version that emits `agent_settled`. Use `/compact [instructions]` to request manual compaction.

## Session mapping

Expand Down
9 changes: 5 additions & 4 deletions docs/04-pi-rpc-integration.md
Original file line number Diff line number Diff line change
Expand Up @@ -68,7 +68,8 @@ Events then stream asynchronously:

```json
{"type":"message_end","message":{...}}
{"type":"agent_end","messages":[...]}
{"type":"agent_end","messages":[...],"willRetry":false}
{"type":"agent_settled"}
```

## Request/Response Handling
Expand Down Expand Up @@ -99,17 +100,17 @@ future resolves
{"type":"prompt","message":"..."}
```

Then it consumes events until `agent_end`.
Then it consumes events until `agent_settled`. `agent_end` closes only a low-level attempt: Pi may compact, retry an overflow, or process queued follow-ups before settling. Before sending a new prompt, the client drains old command events; it never drains events after sending, because Pi can finish before its command response arrives. Each new `agent_start` clears the previous attempt's candidate reply.

The final assistant text is extracted from either:
The final assistant text from the last attempt is extracted from either:

- the last assistant `message_end`, or
- `agent_end.finalText` / `agent_end.final_text`, or
- the `agent_end.messages` array as a compatibility fallback only when no final text was observed earlier.

When `agent_end.messages` is present, `PromptResult.events` omits that message history and records `messagesOmitted`/`messageCount` metadata instead. This prevents gateway callers from retaining a full session snapshot in memory when Telegram only needs the final assistant response.

This avoids requiring token-by-token Telegram streaming for v1.
This avoids requiring token-by-token Telegram streaming for v1. It relies on a Pi version that emits `agent_settled`; older Pi versions that lack this event must be upgraded (the gateway intentionally does not fall back to `agent_end`, which could send an incomplete reply).

## Supported Pi Operations

Expand Down
17 changes: 13 additions & 4 deletions pi_gateway/pi_rpc.py
Original file line number Diff line number Diff line change
Expand Up @@ -214,25 +214,34 @@ async def request(self, payload: dict[str, Any], timeout: float | None = 300) ->
raise PiRpcError(str(response.get("error") or response))
return response

async def events_until_agent_end(self, timeout: float | None = None) -> AsyncIterator[dict[str, Any]]:
async def events_until_agent_settled(self, timeout: float | None = None) -> AsyncIterator[dict[str, Any]]:
while True:
event = await asyncio.wait_for(self._events.get(), timeout=timeout)
if event.get("type") == RPC_ERROR_EVENT:
raise PiRpcError(str(event.get("error") or "pi rpc client failed"))
yield event
if event.get("type") == "agent_end":
if event.get("type") == "agent_settled":
return

async def prompt(self, message: str, *, streaming_behavior: str | None = None) -> PromptResult:
payload: dict[str, Any] = {"type": "prompt", "message": message}
if streaming_behavior:
payload["streamingBehavior"] = streaming_behavior
await self.start()
# Old command events (e.g. a manual compact) have no request ID. Drain
# them before sending the next prompt, never after: a fast Pi response
# may already have queued this prompt's completion by then.
self._clear_events()
await self.request(payload, timeout=60)
events: list[dict[str, Any]] = []
final_text = ""
async for event in self.events_until_agent_end(timeout=None):
async for event in self.events_until_agent_settled(timeout=None):
events.append(event_without_message_history(event))
if event.get("type") == "message_end":
if event.get("type") == "agent_start":
# Retries after an overflow begin another run. Do not send a
# partial answer from the earlier attempt to Telegram.
final_text = ""
elif event.get("type") == "message_end":
msg = event.get("message") or {}
if isinstance(msg, dict) and msg.get("role") == "assistant":
final_text = content_to_text(msg.get("content")) or final_text
Expand Down
90 changes: 90 additions & 0 deletions tests/test_pi_rpc_settled.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,90 @@
import asyncio
import unittest
from unittest.mock import AsyncMock

from pi_gateway.config import PiConfig
from pi_gateway.pi_rpc import RPC_ERROR_EVENT, PiRpcClient, PiRpcError


class PiRpcSettledTest(unittest.IsolatedAsyncioTestCase):
def setUp(self):
self.client = PiRpcClient(PiConfig())
self.client.start = AsyncMock()
self.client.request = AsyncMock(return_value={"success": True})

async def send(self, *events):
for event in events:
await self.client._events.put(event)

async def test_compaction_retry_returns_final_run_only_after_settled(self):
await self.send(
{"type": "agent_start"},
{"type": "message_end", "message": {"role": "assistant", "content": "partial answer"}},
{"type": "agent_end", "willRetry": True, "messages": [{"role": "assistant", "content": "partial answer"}]},
{"type": "compaction_start", "reason": "overflow"},
{"type": "compaction_end"},
{"type": "agent_start"},
{"type": "message_end", "message": {"role": "assistant", "content": [{"type": "text", "text": "final answer"}]}},
{"type": "agent_end", "messages": [{"role": "assistant", "content": "final answer"}]},
{"type": "agent_settled"},
)
# Pi can emit events before it replies to the prompt command. They must
# remain queued when request() returns, not be cleared afterwards.
pending = [self.client._events.get_nowait() for _ in range(self.client._events.qsize())]

async def fast_request(*_args, **_kwargs):
await self.send(*pending)
return {"success": True}

self.client.request.side_effect = fast_request
result = await asyncio.wait_for(self.client.prompt("hi"), timeout=1)
self.assertEqual(result.text, "final answer")
self.assertEqual(result.events[-1]["type"], "agent_settled")
self.assertEqual(len([e for e in result.events if e["type"] == "agent_end"]), 2)
self.assertTrue(all("messages" not in e for e in result.events if e["type"] == "agent_end"))
self.assertTrue(self.client._events.empty())

async def test_stale_completion_from_earlier_command_is_ignored(self):
await self.send({"type": "message_end", "message": {"role": "assistant", "content": "old"}},
{"type": "agent_settled"})

async def new_request(*_args, **_kwargs):
await self.send({"type": "agent_start"},
{"type": "message_end", "message": {"role": "assistant", "content": "new"}},
{"type": "agent_end"}, {"type": "agent_settled"})
return {"success": True}

self.client.request.side_effect = new_request
result = await asyncio.wait_for(self.client.prompt("second"), timeout=1)
self.assertEqual(result.text, "new")
self.assertEqual(len(result.events), 4)

async def test_fallback_text_uses_last_run_not_failed_attempt(self):
async def request(*_args, **_kwargs):
await self.send({"type": "agent_end", "finalText": "partial"},
{"type": "agent_start"},
{"type": "agent_end", "messages": [{"role": "assistant", "content": "recovered"}]},
{"type": "agent_settled"})
return {"success": True}

self.client.request.side_effect = request
result = await asyncio.wait_for(self.client.prompt("overflow"), timeout=1)
self.assertEqual(result.text, "recovered")

async def test_abort_settles_without_assistant_text(self):
async def request(*_args, **_kwargs):
await self.send({"type": "agent_start"}, {"type": "agent_end"}, {"type": "agent_settled"})
return {"success": True}

self.client.request.side_effect = request
result = await asyncio.wait_for(self.client.prompt("abort"), timeout=1)
self.assertEqual(result.text, "")

async def test_reader_failure_is_reported_before_settlement(self):
async def request(*_args, **_kwargs):
await self.send({"type": "agent_end"}, {"type": RPC_ERROR_EVENT, "error": "Pi exited"})
return {"success": True}

self.client.request.side_effect = request
with self.assertRaisesRegex(PiRpcError, "Pi exited"):
await asyncio.wait_for(self.client.prompt("hi"), timeout=1)
Loading