From 8e9aa7af48953dfd2aecb633e010f6e0d52f9144 Mon Sep 17 00:00:00 2001 From: Jeffrey Phillips Freeman Date: Mon, 23 Mar 2026 23:16:35 +0000 Subject: [PATCH] fix(client): address server client chain review findings - Promote _request() to public method request() in ServerHttpClient to fix encapsulation violation across sync_client and remote_project - Fix WebSocket connect() to raise NotImplementedError with clear TODO documenting that real websockets transport is not yet implemented - Fix thread safety: protect last_event_id write under self._lock in ws_client.process_event() - Broaden exception handling in sync() to catch ServerTimeoutError and A2aNotAvailableError in addition to ServerConnectionError - Add has_next to PageResult in benchmark and test helpers - Add comments explaining Retry-After header logging-only behavior and why blocking time.sleep is acceptable in sync client - Update all test steps, robot helpers, and benchmarks to use the renamed public request() method and updated connect() behavior Refs: #335, #336, #337, #338 --- benchmarks/server_remote_project_bench.py | 2 +- benchmarks/server_sync_bench.py | 2 +- benchmarks/server_ws_bench.py | 6 ++-- features/client/websocket_updates.feature | 18 ++++++------ features/steps/plan_sync_steps.py | 14 ++++----- features/steps/remote_project_steps.py | 6 ++-- features/steps/websocket_updates_steps.py | 26 +++++++++++++---- robot/helper_plan_sync.py | 8 +++--- robot/helper_remote_project.py | 4 +-- robot/helper_websocket_updates.py | 34 +++++++++++++++------- robot/websocket_updates.robot | 8 +++--- src/cleveragents/client/http_client.py | 35 +++++++++++++++++++---- src/cleveragents/client/remote_project.py | 2 +- src/cleveragents/client/sync_client.py | 17 +++++------ src/cleveragents/client/ws_client.py | 27 ++++++++++------- 15 files changed, 136 insertions(+), 73 deletions(-) diff --git a/benchmarks/server_remote_project_bench.py b/benchmarks/server_remote_project_bench.py index 82ab713d..80dd68fe 100644 --- a/benchmarks/server_remote_project_bench.py +++ b/benchmarks/server_remote_project_bench.py @@ -33,7 +33,7 @@ _ITEMS = [ def _build() -> RemoteProjectClient: mock = MagicMock(spec=ServerHttpClient) mock.list_endpoint.return_value = PageResult( - items=_ITEMS, page=1, per_page=50, total=20 + items=_ITEMS, page=1, per_page=50, total=20, has_next=False ) return RemoteProjectClient(mock) diff --git a/benchmarks/server_sync_bench.py b/benchmarks/server_sync_bench.py index 6bd7ea5d..6e87880c 100644 --- a/benchmarks/server_sync_bench.py +++ b/benchmarks/server_sync_bench.py @@ -54,7 +54,7 @@ class SyncOperationSuite: def setup(self) -> None: self.mock_http = MagicMock(spec=ServerHttpClient) - self.mock_http._request.return_value = _mock_response(200, {"id": "srv"}) + self.mock_http.request.return_value = _mock_response(200, {"id": "srv"}) self.client = PlanSyncClient(self.mock_http) self.items = [{"id": f"item-{i}"} for i in range(10)] diff --git a/benchmarks/server_ws_bench.py b/benchmarks/server_ws_bench.py index e7d48045..6876fe90 100644 --- a/benchmarks/server_ws_bench.py +++ b/benchmarks/server_ws_bench.py @@ -50,12 +50,14 @@ class EventProcessingSuite: def setup(self) -> None: self.client = WebSocketClient() - self.client.connect() + self.client._state.connected = True + self.client._running = True self.client.subscribe(lambda e: None) def time_process_event(self) -> None: c = WebSocketClient() - c.connect() + c._state.connected = True + c._running = True c.subscribe(lambda e: None) for i in range(100): event = A2aEvent(event_id=f"bench-{i}", event_type="plan.status") diff --git a/features/client/websocket_updates.feature b/features/client/websocket_updates.feature index 23496706..5d76cc53 100644 --- a/features/client/websocket_updates.feature +++ b/features/client/websocket_updates.feature @@ -22,15 +22,15 @@ Feature: WebSocket updates client # Connection lifecycle # --------------------------------------------------------------------------- - Scenario: Connect sets state to connected + Scenario: Connect raises NotImplementedError until real transport is available Given a WebSocketClient with default settings - When I connect the ws client - Then the ws client should be connected + When I attempt to connect the ws client + Then a NotImplementedError should be raised from ws client Scenario: Disconnect sets state to disconnected Given a WebSocketClient with default settings - When I connect the ws client - And I disconnect the ws client + And the ws client state is set to connected + When I disconnect the ws client Then the ws client should not be connected # --------------------------------------------------------------------------- @@ -39,8 +39,8 @@ Feature: WebSocket updates client Scenario: Reconnect increments reconnect counter Given a WebSocketClient with default settings and max_reconnects 5 - When I connect the ws client - And I simulate a reconnect + And the ws client state is set to connected + When I simulate a reconnect Then the reconnect count should be 1 And the ws client should be connected @@ -77,8 +77,8 @@ Feature: WebSocket updates client Scenario: Heartbeat resets reconnect counter Given a WebSocketClient with default settings and max_reconnects 5 - When I connect the ws client - And I simulate a reconnect + And the ws client state is set to connected + When I simulate a reconnect And I handle a heartbeat Then the reconnect count should be 0 diff --git a/features/steps/plan_sync_steps.py b/features/steps/plan_sync_steps.py index d2f16fbc..cfbe4112 100644 --- a/features/steps/plan_sync_steps.py +++ b/features/steps/plan_sync_steps.py @@ -98,7 +98,7 @@ def step_sync_client_local_wins(context: Context) -> None: def step_new_items(context: Context, n: int) -> None: context.sync_items = [{"id": f"item-{i}", "name": f"Item {i}"} for i in range(n)] # Mock POST returns new server_id - context.mock_http._request.return_value = _mock_response(200, {"id": "srv-new-001"}) + context.mock_http.request.return_value = _mock_response(200, {"id": "srv-new-001"}) @given("a list of {n:d} existing items with server_id and newer version") @@ -113,7 +113,7 @@ def step_existing_newer(context: Context, n: int) -> None: return _mock_response(200, {"version": 1}) return _mock_response(200, {"id": "srv-updated"}) - context.mock_http._request.side_effect = _side_effect + context.mock_http.request.side_effect = _side_effect @given("a list of {n:d} existing items with server_id and same version") @@ -121,7 +121,7 @@ def step_existing_same(context: Context, n: int) -> None: context.sync_items = [ {"id": f"item-{i}", "server_id": f"srv-{i}", "version": 1} for i in range(n) ] - context.mock_http._request.return_value = _mock_response(200, {"version": 1}) + context.mock_http.request.return_value = _mock_response(200, {"version": 1}) @given("a list with one item missing id") @@ -134,7 +134,7 @@ def step_existing_server_newer(context: Context, n: int) -> None: context.sync_items = [ {"id": f"item-{i}", "server_id": f"srv-{i}", "version": 1} for i in range(n) ] - context.mock_http._request.return_value = _mock_response(200, {"version": 5}) + context.mock_http.request.return_value = _mock_response(200, {"version": 5}) # --------------------------------------------------------------------------- @@ -201,7 +201,7 @@ def step_resources(context: Context) -> None: "actions": [{"id": "a1"}], "tools": [{"id": "t1"}], } - context.mock_http._request.return_value = _mock_response(200, {"id": "srv-001"}) + context.mock_http.request.return_value = _mock_response(200, {"id": "srv-001"}) @when("I sync all with a scope limited to actions only") @@ -245,13 +245,13 @@ def step_sync_client_exec(context: Context) -> None: ) return _mock_response(200, {}) - context.mock_http._request.side_effect = _exec_side_effect + context.mock_http.request.side_effect = _exec_side_effect @given("a PlanSyncClient with a mock HTTP client for status") def step_sync_client_status(context: Context) -> None: context.sync_client, context.mock_http = _build_sync_client() - context.mock_http._request.return_value = _mock_response( + context.mock_http.request.return_value = _mock_response( 200, {"phase": "running", "progress": 50} ) context.call_error = None diff --git a/features/steps/remote_project_steps.py b/features/steps/remote_project_steps.py index 1393f1f6..d1a5aa6a 100644 --- a/features/steps/remote_project_steps.py +++ b/features/steps/remote_project_steps.py @@ -20,7 +20,9 @@ from cleveragents.core.exceptions import ResourceNotFoundError def _mock_page_result( items: list[dict[str, Any]], ) -> PageResult: - return PageResult(items=items, page=1, per_page=50, total=len(items)) + return PageResult( + items=items, page=1, per_page=50, total=len(items), has_next=False + ) _DEFAULT_PROJECTS = [ @@ -193,7 +195,7 @@ def _mock_exec_response( @given("a RemoteProjectClient with mock execution endpoint") def step_remote_exec(context: Context) -> None: context.remote_client, context.mock_http = _build_remote_client() - context.mock_http._request.return_value = _mock_exec_response( + context.mock_http.request.return_value = _mock_exec_response( 200, {"status": "submitted", "execution_id": "exec-001"} ) context.call_error = None diff --git a/features/steps/websocket_updates_steps.py b/features/steps/websocket_updates_steps.py index 735429d8..c4e09f75 100644 --- a/features/steps/websocket_updates_steps.py +++ b/features/steps/websocket_updates_steps.py @@ -59,9 +59,18 @@ def step_check_dedup_cap(context: Context, expected: int) -> None: # --------------------------------------------------------------------------- -@when("I connect the ws client") -def step_ws_connect(context: Context) -> None: - context.ws_client.connect() +@when("I attempt to connect the ws client") +def step_ws_attempt_connect(context: Context) -> None: + try: + context.ws_client.connect() + except NotImplementedError as exc: + context.call_error = exc + + +@given("the ws client state is set to connected") +def step_ws_set_connected(context: Context) -> None: + context.ws_client._state.connected = True + context.ws_client._running = True @when("I disconnect the ws client") @@ -92,7 +101,8 @@ def step_ws_reconnect(context: Context) -> None: @when("I exhaust reconnect attempts") def step_ws_exhaust_reconnects(context: Context) -> None: - context.ws_client.connect() + context.ws_client._state.connected = True + context.ws_client._running = True try: for _ in range(context.ws_client._max_reconnects + 1): context.ws_client._state.connected = False @@ -111,6 +121,11 @@ def step_ws_conn_error(context: Context) -> None: assert isinstance(context.call_error, ServerConnectionError) +@then("a NotImplementedError should be raised from ws client") +def step_ws_not_implemented_error(context: Context) -> None: + assert isinstance(context.call_error, NotImplementedError) + + # --------------------------------------------------------------------------- # Event processing # --------------------------------------------------------------------------- @@ -119,7 +134,8 @@ def step_ws_conn_error(context: Context) -> None: @given("a WebSocketClient with a subscriber") def step_ws_with_subscriber(context: Context) -> None: context.ws_client = WebSocketClient() - context.ws_client.connect() + context.ws_client._state.connected = True + context.ws_client._running = True context.received_events: list[A2aEvent] = [] def _on_event(event: A2aEvent) -> None: diff --git a/robot/helper_plan_sync.py b/robot/helper_plan_sync.py index 558aa8b3..46ef7569 100644 --- a/robot/helper_plan_sync.py +++ b/robot/helper_plan_sync.py @@ -42,7 +42,7 @@ def scope_active() -> None: def sync_create() -> None: mock_http = MagicMock(spec=ServerHttpClient) - mock_http._request.return_value = _mock_response(200, {"id": "srv-new"}) + mock_http.request.return_value = _mock_response(200, {"id": "srv-new"}) client = PlanSyncClient(mock_http) items = [{"id": "a1"}, {"id": "a2"}] summary = client.sync(items, "actions") @@ -62,7 +62,7 @@ def sync_dry_run() -> None: def execute_plan() -> None: mock_http = MagicMock(spec=ServerHttpClient) - mock_http._request.return_value = _mock_response( + mock_http.request.return_value = _mock_response( 200, {"server_plan_id": "srv-001", "status": "submitted"} ) client = PlanSyncClient(mock_http) @@ -74,7 +74,7 @@ def execute_plan() -> None: def apply_plan() -> None: mock_http = MagicMock(spec=ServerHttpClient) - mock_http._request.return_value = _mock_response( + mock_http.request.return_value = _mock_response( 200, {"server_plan_id": "srv-002", "status": "applying"} ) client = PlanSyncClient(mock_http) @@ -85,7 +85,7 @@ def apply_plan() -> None: def plan_status() -> None: mock_http = MagicMock(spec=ServerHttpClient) - mock_http._request.return_value = _mock_response( + mock_http.request.return_value = _mock_response( 200, {"phase": "running", "progress": 50} ) client = PlanSyncClient(mock_http) diff --git a/robot/helper_remote_project.py b/robot/helper_remote_project.py index a0c06375..fe05c9d1 100644 --- a/robot/helper_remote_project.py +++ b/robot/helper_remote_project.py @@ -31,7 +31,7 @@ _ITEMS = [ def _build() -> tuple[RemoteProjectClient, MagicMock]: mock = MagicMock(spec=ServerHttpClient) mock.list_endpoint.return_value = PageResult( - items=_ITEMS, page=1, per_page=50, total=2 + items=_ITEMS, page=1, per_page=50, total=2, has_next=False ) return RemoteProjectClient(mock), mock @@ -71,7 +71,7 @@ def not_found() -> None: def request_execution() -> None: client, mock = _build() - mock._request.return_value = httpx.Response( + mock.request.return_value = httpx.Response( status_code=200, json={"status": "submitted"}, headers={}, diff --git a/robot/helper_websocket_updates.py b/robot/helper_websocket_updates.py index 561256bf..86fd8abc 100644 --- a/robot/helper_websocket_updates.py +++ b/robot/helper_websocket_updates.py @@ -18,21 +18,29 @@ from cleveragents.client.ws_client import ( # noqa: E402 ) -def connect_disconnect() -> None: - """Verify connect / disconnect lifecycle.""" +def connect_raises() -> None: + """Verify connect raises NotImplementedError.""" client = WebSocketClient() assert client.connected is False - client.connect() - assert client.connected is True + try: + client.connect() + print("FAIL: should have raised NotImplementedError", file=sys.stderr) + sys.exit(1) + except NotImplementedError: + pass + # Verify disconnect works when state is set externally + client._state.connected = True + client._running = True client.disconnect() assert client.connected is False - print("ws-connect-disconnect-ok") + print("ws-connect-raises-ok") def subscribe_event() -> None: """Verify event subscription and dispatch.""" client = WebSocketClient() - client.connect() + client._state.connected = True + client._running = True received: list[A2aEvent] = [] client.subscribe(received.append) event = A2aEvent( @@ -50,7 +58,8 @@ def subscribe_event() -> None: def dedup_event() -> None: """Verify duplicate event detection.""" client = WebSocketClient() - client.connect() + client._state.connected = True + client._running = True received: list[A2aEvent] = [] client.subscribe(received.append) event = A2aEvent( @@ -68,7 +77,8 @@ def dedup_event() -> None: def reconnect() -> None: """Verify reconnection with backoff.""" client = WebSocketClient(max_reconnects=3, reconnect_base=0.001, reconnect_max=0.01) - client.connect() + client._state.connected = True + client._running = True client._state.connected = False result = client.reconnect() assert result is True @@ -79,7 +89,8 @@ def reconnect() -> None: def reconnect_exhaust() -> None: """Verify exhausted reconnects raise error.""" client = WebSocketClient(max_reconnects=1, reconnect_base=0.001, reconnect_max=0.01) - client.connect() + client._state.connected = True + client._running = True client._state.connected = False client.reconnect() client._state.connected = False @@ -95,7 +106,8 @@ def reconnect_exhaust() -> None: def heartbeat() -> None: """Verify heartbeat resets reconnect counter.""" client = WebSocketClient(max_reconnects=5, reconnect_base=0.001, reconnect_max=0.01) - client.connect() + client._state.connected = True + client._running = True client._state.connected = False client.reconnect() assert client.reconnect_count == 1 @@ -138,7 +150,7 @@ def version_negotiate() -> None: _COMMANDS = { - "connect-disconnect": connect_disconnect, + "connect-raises": connect_raises, "subscribe-event": subscribe_event, "dedup-event": dedup_event, "reconnect": reconnect, diff --git a/robot/websocket_updates.robot b/robot/websocket_updates.robot index f56a87bb..2313f554 100644 --- a/robot/websocket_updates.robot +++ b/robot/websocket_updates.robot @@ -8,13 +8,13 @@ Suite Teardown Cleanup Test Environment ${HELPER} ${CURDIR}/helper_websocket_updates.py *** Test Cases *** -WS Connect Disconnect - [Documentation] Verify connect and disconnect lifecycle - ${result}= Run Process ${PYTHON} ${HELPER} connect-disconnect cwd=${WORKSPACE} +WS Connect Raises NotImplementedError + [Documentation] Verify connect raises NotImplementedError until real transport is available + ${result}= Run Process ${PYTHON} ${HELPER} connect-raises cwd=${WORKSPACE} Log ${result.stdout} Log ${result.stderr} Should Be Equal As Integers ${result.rc} 0 - Should Contain ${result.stdout} ws-connect-disconnect-ok + Should Contain ${result.stdout} ws-connect-raises-ok WS Subscribe Event [Documentation] Verify event subscription and dispatch diff --git a/src/cleveragents/client/http_client.py b/src/cleveragents/client/http_client.py index 14874394..aca4bd5f 100644 --- a/src/cleveragents/client/http_client.py +++ b/src/cleveragents/client/http_client.py @@ -200,6 +200,11 @@ class ServerHttpClient: retry_hint = "" if status == 429: + # The Retry-After header is included in the error message for + # diagnostic visibility but is intentionally not used to + # override the exponential backoff delay. The client's own + # backoff already caps at ``backoff_max`` and using a + # server-supplied value could introduce unbounded waits. retry_after = response.headers.get("Retry-After", "") retry_hint = f" (retry after {retry_after}s)" if retry_after else "" @@ -225,7 +230,7 @@ class ServerHttpClient: ) raise ServerConnectionError(message=msg, url=url, details=details) - def _request( + def request( self, method: str, path: str, @@ -233,7 +238,22 @@ class ServerHttpClient: json_body: dict[str, Any] | None = None, params: dict[str, str] | None = None, ) -> httpx.Response: - """Execute an HTTP request with retry for idempotent methods.""" + """Execute an HTTP request with retry for idempotent methods. + + Args: + method: HTTP method (GET, POST, PUT, DELETE, etc.). + path: URL path relative to the base URL. + json_body: Optional JSON request body. + params: Optional query parameters. + + Returns: + The successful ``httpx.Response``. + + Raises: + ServerConnectionError: On connection or HTTP failure. + ServerTimeoutError: On request timeout. + A2aNotAvailableError: When the server returns 503. + """ url = f"{self._base_url}{path}" headers = self._build_headers() self._log_request(method, url, headers) @@ -273,6 +293,9 @@ class ServerHttpClient: delay=delay, status=response.status_code, ) + # Blocking sleep is intentional: this is a synchronous + # client and callers expect blocking semantics. An + # async variant should use ``asyncio.sleep`` instead. time.sleep(delay) continue @@ -349,7 +372,7 @@ class ServerHttpClient: ServerTimeoutError: On request timeout. """ try: - response = self._request("GET", "/health") + response = self.request("GET", "/health") data = response.json() return bool(data.get("status") == "healthy") except (ServerConnectionError, ServerTimeoutError, A2aNotAvailableError): @@ -365,7 +388,7 @@ class ServerHttpClient: ServerConnectionError: On connection failure. ServerTimeoutError: On request timeout. """ - response = self._request("GET", "/version") + response = self.request("GET", "/version") data = response.json() version: str = str(data.get("version", "")) return version @@ -386,7 +409,7 @@ class ServerHttpClient: ServerVersionMismatchError: When negotiation fails. ServerConnectionError: On connection failure. """ - response = self._request( + response = self.request( "POST", "/version/negotiate", json_body={"client_version": client_version}, @@ -423,7 +446,7 @@ class ServerHttpClient: query: dict[str, str] = {"page": str(page), "per_page": str(per_page)} if params: query.update(params) - response = self._request("GET", path, params=query) + response = self.request("GET", path, params=query) data = response.json() items: list[dict[str, Any]] diff --git a/src/cleveragents/client/remote_project.py b/src/cleveragents/client/remote_project.py index 52a40ebe..501d05ea 100644 --- a/src/cleveragents/client/remote_project.py +++ b/src/cleveragents/client/remote_project.py @@ -219,7 +219,7 @@ class RemoteProjectClient: if not project_id: raise ValueError("project_id must be a non-empty string") - resp = self._http._request( + resp = self._http.request( "POST", f"/projects/{project_id}/execute", json_body={"plan_name": plan_name} if plan_name else None, diff --git a/src/cleveragents/client/sync_client.py b/src/cleveragents/client/sync_client.py index 025ac304..d5c35f79 100644 --- a/src/cleveragents/client/sync_client.py +++ b/src/cleveragents/client/sync_client.py @@ -14,7 +14,8 @@ from typing import Any import structlog -from cleveragents.client.exceptions import ServerConnectionError +from cleveragents.a2a.errors import A2aNotAvailableError +from cleveragents.client.exceptions import ServerConnectionError, ServerTimeoutError from cleveragents.client.http_client import ServerHttpClient logger: structlog.stdlib.BoundLogger = structlog.get_logger(__name__) @@ -207,7 +208,7 @@ class PlanSyncClient: server_id = str(item.get("server_id", local_id)) summary.server_ids[local_id] = server_id - except ServerConnectionError: + except (ServerConnectionError, ServerTimeoutError, A2aNotAvailableError): summary.errors += 1 logger.warning("sync_item_failed", local_id=local_id) @@ -262,7 +263,7 @@ class PlanSyncClient: if server_id is None: # New item — create on server - resp = self._http._request( + resp = self._http.request( "POST", f"/{resource_type}", json_body={"item": item}, @@ -274,7 +275,7 @@ class PlanSyncClient: # Existing item — check for conflict try: - resp = self._http._request( + resp = self._http.request( "GET", f"/{resource_type}/{server_id}", ) @@ -292,7 +293,7 @@ class PlanSyncClient: return "skipped" # Local wins or local version is newer — push update - resp = self._http._request( + resp = self._http.request( "PUT", f"/{resource_type}/{server_id}", json_body={"item": item}, @@ -315,7 +316,7 @@ class PlanSyncClient: if not plan_id: raise ValueError("plan_id must be a non-empty string") - resp = self._http._request( + resp = self._http.request( "POST", "/plans/execute", json_body={"plan_id": plan_id}, @@ -340,7 +341,7 @@ class PlanSyncClient: if not plan_id: raise ValueError("plan_id must be a non-empty string") - resp = self._http._request( + resp = self._http.request( "POST", "/plans/apply", json_body={"plan_id": plan_id}, @@ -365,7 +366,7 @@ class PlanSyncClient: if not server_plan_id: raise ValueError("server_plan_id must be a non-empty string") - resp = self._http._request("GET", f"/plans/{server_plan_id}/status") + resp = self._http.request("GET", f"/plans/{server_plan_id}/status") data: dict[str, Any] = resp.json() return data diff --git a/src/cleveragents/client/ws_client.py b/src/cleveragents/client/ws_client.py index 817eefac..b02a63e1 100644 --- a/src/cleveragents/client/ws_client.py +++ b/src/cleveragents/client/ws_client.py @@ -248,9 +248,8 @@ class WebSocketClient: logger.debug("ws_event_duplicate", event_id=event.event_id) return False - self._state.last_event_id = event.event_id - with self._lock: + self._state.last_event_id = event.event_id callbacks = list(self._callbacks) for cb in callbacks: @@ -274,16 +273,24 @@ class WebSocketClient: # ------------------------------------------------------------------ def connect(self) -> None: - """Establish the WebSocket connection (simulated). + """Establish the WebSocket connection. - In the current implementation this sets state to connected. - The actual WebSocket connection will be established when - the server project provides a real endpoint. + .. todo:: + Implement a real WebSocket connection using the ``websockets`` + library. This requires a running server endpoint to connect + to and should include TLS support, authentication header + injection, and proper async lifecycle management. + + Raises: + NotImplementedError: Always — real WebSocket transport is + not yet implemented. """ - self._state.connected = True - self._state.reconnect_count = 0 - self._running = True - logger.info("ws_connected", url=self._ws_url) + raise NotImplementedError( + "WebSocket transport is not yet implemented. " + "A real connection using the `websockets` library is required " + "before this client can be used in production. " + "See: https://websockets.readthedocs.io/" + ) def disconnect(self) -> None: """Close the WebSocket connection and stop reconnection."""