Skip to content
Draft
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: 2 additions & 0 deletions src/band/client/streaming/client.py
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,7 @@
PhoenixChannelsProtocolVersion,
PHXChannelsClient,
)
from phoenix_channels_python_client.client_types import ReconnectPolicy
from phoenix_channels_python_client.exceptions import PHXConnectionError
from phoenix_channels_python_client.phx_messages import PHXMessage
from pydantic import Field
Expand Down Expand Up @@ -450,6 +451,7 @@ def _build_phx_client(self) -> PHXChannelsClient:
self.api_key,
protocol_version=PhoenixChannelsProtocolVersion.V2,
auto_reconnect=False,
reconnect_policy=ReconnectPolicy(reconnect_on_normal_close=True),
heartbeat_interval_s=self._watchdog.policy.heartbeat_interval_s,
on_reconnect=self._handle_reconnect,
on_disconnect=self._on_disconnect,
Expand Down
49 changes: 49 additions & 0 deletions tests/platform/test_link_real_wire.py
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@

from band.core.exceptions import AgentAlreadyRunningError
from band.core.types import ConflictPolicy
from band.platform.event import WebSocketDisconnectedEvent
from band.platform.link import BandLink
from band.testing import FakePhoenixServer, JoinOutcome, fake_phoenix_server
from tests.conftest import spy_on_reconciliation_drain
Expand Down Expand Up @@ -191,3 +192,51 @@ async def test_refused_duplicate_raises_agent_already_running() -> None:
ConflictPolicy.REJECT,
ConflictPolicy.REJECT,
]


@pytest.mark.parametrize("code", [1000, 1001])
async def test_server_clean_close_without_supersede_reconnects(code: int) -> None:
async with fake_phoenix_server() as server:
link = make_link(server.url)
await link.connect()
reconnect_handled = spy_on_reconciliation_drain(link)

await server.close_connection(code=code)
await asyncio.wait_for(reconnect_handled.wait(), timeout=5.0)

assert server.connection_count == 2
assert "agent_control:agent-123" in server.joined_topics
assert link.last_disconnect_reason is None


async def test_terminal_supersede_stays_down_after_clean_close() -> None:
async with fake_phoenix_server() as server:
link = make_link(server.url)
await link.connect()
await server.push(
"agent_control:agent-123",
"supersede",
{
"reason": "session.already_connected",
"message": "Superseded by a newer session",
"retryable": False,
"correlation_id": "evict-123",
},
)
event = await asyncio.wait_for(link.__anext__(), timeout=2.0)
assert isinstance(event, WebSocketDisconnectedEvent)
assert event.payload.reason == "session.already_connected"

await server.close_connection(code=1000)
await asyncio.sleep(1.0)
assert server.connection_count == 1
assert link.last_disconnect_reason is not None


async def test_local_disconnect_stays_down_with_clean_close_retries_enabled() -> None:
async with fake_phoenix_server() as server:
link = make_link(server.url)
await link.connect()
await link.disconnect()
await asyncio.sleep(1.0)
assert server.connection_count == 1
Loading