From 0a400b5482d6fd1fed6fb9365da718308ff18d4f Mon Sep 17 00:00:00 2001 From: Federico Lois Date: Fri, 19 Sep 2025 12:38:02 -0300 Subject: [PATCH] RDBC-948: Ensure topology handling matches C# Client. Handles network errors during connection failover more reliably by catching both requests-level and OS-level socket exceptions. This mirrors the behavior of the C# client and prevents connection failures due to DNS resolution issues or server unavailability. Also, refactors the failover logic to improve retry behavior on failed requests, ensuring that the next preferred node is chosen and retried, while avoiding infinite loops by checking for already failed nodes. Notifies listeners about failed requests with full details, and broadcasts commands where applicable to improve efficiency. Finally, ensures that only member nodes without failures are chosen first. Adds tests to verify failover scenarios with invalid DNS configurations and unreachable endpoints. --- ravendb/http/request_executor.py | 93 +++++++++-------- ravendb/http/topology.py | 3 +- ravendb/tests/issue_tests/test_RDBC_948.py | 113 +++++++++++++++++++++ 3 files changed, 168 insertions(+), 41 deletions(-) create mode 100644 ravendb/tests/issue_tests/test_RDBC_948.py diff --git a/ravendb/http/request_executor.py b/ravendb/http/request_executor.py index 0658fc84..2c8314dc 100644 --- a/ravendb/http/request_executor.py +++ b/ravendb/http/request_executor.py @@ -37,7 +37,7 @@ from http import HTTPStatus -from typing import TYPE_CHECKING, List, Dict, Tuple +from typing import TYPE_CHECKING, List, Dict, Tuple, Optional if TYPE_CHECKING: from ravendb.documents.session import SessionInfo @@ -633,7 +633,13 @@ def _send_request_to_server( self._throw_failed_to_contact_all_nodes(command, request) return None - except IOError as e: + except (requests.RequestException, OSError) as e: + # RDBC-948: https://issues.hibernatingrhinos.com/issue/RDBC-948/Python-client-connection-failover-breaks-with-unknown-DNS-name-or-server-is-down. + # Handle failover on network errors from both requests and the OS: + # - RequestException covers requests' network stack (connect, TLS, proxies, etc.). + # - OSError covers socket-level issues like DNS getaddrinfo on some platforms. + # Different OS/resolvers surface the same fault differently; catching both mirrors the C# client + # (HttpRequestException/SocketException) and makes failover reliable. if not should_retry: raise @@ -805,7 +811,8 @@ def _throw_failed_to_contact_all_nodes(self, command: RavenCommand, request: req ) if len(command.failed_nodes) == 1: - raise command.failed_nodes.popitem() + # raise the single recorded exception + raise next(iter(command.failed_nodes.values())) message = ( f"Tried to send {command._result_class.__name__} request via {request.method}" @@ -1159,39 +1166,45 @@ def __handle_server_down( if command.failed_nodes is None: command.failed_nodes = {} - return ( - False # todo: command.failed_nodes[chosen_node] = self.__read_exception_from_server(request, response, e) - ) + # record the failure for this node + if not command.is_failed_with_node(chosen_node): + command.failed_nodes[chosen_node] = self.__read_exception_from_server(request, response, e) + # If the node is not part of the topology, we can't failover using selector. if node_index is None: - # We executed request over a node not in the topology. This means no failover... return False + # If we don't have a selector yet, we also cannot failover. if self._node_selector is None: - # todo: spawnHealthChecks(chosenNode, nodeIndex) return False - # As the server is down, we discard the server version to ensure we update when it goes up. + # As the server is down, discard server version to ensure it updates when back up. chosen_node.discard_server_version() + # Mark the node as failed to move selection forward self._node_selector.on_failed_request(node_index) + # For broadcastable commands, attempt to broadcast instead of single-node retry. if self.should_broadcast(command): command.result = self.__broadcast(command, session_info) return True - # todo: self.spawn_health_checks(chosen_node, node_index) - + # Choose the next preferred node and retry index_node_and_etag = self._node_selector.get_preferred_node_with_topology() + + # If topology changed since we started, clear failed nodes record to allow retries if command.failover_topology_etag != self.topology_etag: command.failed_nodes.clear() command.failover_topology_etag = self.topology_etag + # Avoid infinite loop if the next node is already marked as failed if index_node_and_etag.current_node in command.failed_nodes: return False - self.__on_failed_request_invoke(url, e, request, response) + # Notify listeners about the failed request with full details + self.__on_failed_request_invoke_details(url, e, request, response) + # Retry the command on the next node self.execute( index_node_and_etag.current_node, index_node_and_etag.current_index, command, should_retry, session_info ) @@ -1199,41 +1212,41 @@ def __handle_server_down( return True @staticmethod - def __read_exception_from_server(request: requests.Request, response: requests.Response, e: Exception) -> Exception: - if response and response.content: - response_json = None + def __read_exception_from_server( + request: requests.Request, response: requests.Response, e: Optional[Exception] + ) -> Exception: + # Prefer server-provided error when available + if response is not None and response.content: + raw = None try: - response_json = response.content.decode("utf-8") + raw = response.content.decode("utf-8") - # todo: change this bs - def exception_schema_decoder(dictionary: dict) -> ExceptionDispatcher.ExceptionSchema: + def _decode(d: dict) -> ExceptionDispatcher.ExceptionSchema: return ExceptionDispatcher.ExceptionSchema( - dictionary.get("url"), - dictionary.get("class"), - dictionary.get("message"), - dictionary.get("error"), + d.get("url"), d.get("class"), d.get("message"), d.get("error") ) - return ExceptionDispatcher.get( - json.loads(response_json, object_hook=exception_schema_decoder), response.status_code, e - ) - except: - exception_schema = ExceptionDispatcher.ExceptionSchema( - request.url, + return ExceptionDispatcher.get(json.loads(raw, object_hook=_decode), response.status_code, e) + except Exception: + schema = ExceptionDispatcher.ExceptionSchema( + request.url if request else "", "Unparsable Server Response", - "Get unrecognized response from the server", - response_json, + "Unrecognized response from server", + raw, ) - - return ExceptionDispatcher.get(exception_schema, response.status_code, e) - - exception_schema = ExceptionDispatcher.ExceptionSchema( - request.url, - e.__class__.__qualname__, - e.args[0], - f"An exception occurred while contacting {request.url}.{os.linesep}{str(e)}", - ) - return ExceptionDispatcher.get(exception_schema, HTTPStatus.SERVICE_UNAVAILABLE, e) + return ExceptionDispatcher.get(schema, response.status_code, e) + + # Fallback when we have no usable response body + url = request.url if request else "" + cls = type(e).__name__ if e else "RequestFailed" + msg = str(e) if e else "Request failed" + details = f"An exception occurred while contacting {url}." + if e: + details += f"{os.linesep}{msg}" + + schema = ExceptionDispatcher.ExceptionSchema(url, cls, msg, details) + status = response.status_code if response else HTTPStatus.SERVICE_UNAVAILABLE + return ExceptionDispatcher.get(schema, status, e or RuntimeError(msg)) class IndexAndResponse: def __init__(self, index: int, response: requests.Response): diff --git a/ravendb/http/topology.py b/ravendb/http/topology.py index fd2dd78c..f555ed5c 100644 --- a/ravendb/http/topology.py +++ b/ravendb/http/topology.py @@ -149,7 +149,8 @@ def get_preferred_node_internal(cls, state: NodeSelector.__NodeSelectorState) -> server_nodes = state.nodes length = min(len(server_nodes), len(state_failures)) for i in range(length): - if state_failures[0] == 0: + # pick the first node without failures + if state_failures[i] == 0 and server_nodes[i].server_role == ServerNode.Role.MEMBER: return CurrentIndexAndNode(i, server_nodes[i]) return cls.unlikely_everyone_faulted_choice(state) diff --git a/ravendb/tests/issue_tests/test_RDBC_948.py b/ravendb/tests/issue_tests/test_RDBC_948.py new file mode 100644 index 00000000..a8b95185 --- /dev/null +++ b/ravendb/tests/issue_tests/test_RDBC_948.py @@ -0,0 +1,113 @@ +from threading import Event +from ravendb.documents.store.definition import DocumentStore +from ravendb.documents.subscriptions.options import SubscriptionWorkerOptions +from ravendb.exceptions.exceptions import AllTopologyNodesDownException +from ravendb.infrastructure.entities import User +from ravendb.serverwide.operations.common import GetDatabaseRecordOperation +from ravendb.tests.test_base import TestBase +from ravendb.http.topology import UpdateTopologyParameters +from ravendb.http.server_node import ServerNode + + +class TestRDBC948(TestBase): + def setUp(self): + super().setUp() + + def test_failover_with_invalid_dns_in_urls(self): + # One invalid DNS hostname and one valid server URL (from the embedded test server) + invalid_host = "http://thisnamedoesnotexist:8080" + valid_url = self.store.urls[0] + + with DocumentStore(urls=[invalid_host, valid_url], database=self.store.database) as store2: + store2.conventions.disable_topology_updates = False + store2.initialize() + + # Should succeed by failing over to the valid URL + with store2.open_session() as session: + session.store({"Name": "John"}, "users/1") + session.save_changes() + + # Verify we can read it back (continues using the healthy node) + with store2.open_session() as session: + doc = session.load("users/1") + self.assertIsNotNone(doc) + self.assertEqual(doc.get("Name"), "John") + + def test_all_nodes_down_throws(self): + # Two unreachable endpoints: invalid DNS and a closed localhost port + urls = [ + "http://thisnamedoesnotexist:8080", + "http://127.0.0.1:1234", + ] + + with DocumentStore(urls=urls, database=self.store.database) as store2: + store2.conventions.disable_topology_updates = False + store2.initialize() + + with self.assertRaises(AllTopologyNodesDownException): + with store2.open_session() as session: + session.load("users/does-not-matter") + + def test_maintenance_operation_failover_with_invalid_dns(self): + invalid_host = "http://thisnamedoesnotexist:8080" + valid_url = self.store.urls[0] + database = self.store.database + + with DocumentStore(urls=[invalid_host, valid_url], database=database) as store2: + store2.conventions.disable_topology_updates = False + store2.initialize() + + # Perform maintenance call, should succeed by failing over + record = store2.maintenance.server.send(GetDatabaseRecordOperation(database)) + self.assertIsNotNone(record) + + def test_request_executor_failover_with_invalid_dns(self): + invalid_host = "http://thisnamedoesnotexist:8080" + valid_url = self.store.urls[0] + + with DocumentStore(urls=[invalid_host, valid_url], database=self.store.database) as store2: + store2.conventions.disable_topology_updates = False + store2.initialize() + + # Explicitly refresh topology like C# tests via UpdateTopologyAsync. + # This avoids racing the background first-topology-update and ensures the selector is initialized. + req_ex = store2.get_request_executor() + params = UpdateTopologyParameters(ServerNode(valid_url, store2.database)) + params.timeout_in_ms = 5 + params.debug_tag = "test-init" + req_ex.update_topology_async(params).result() + # Now URL should reflect the healthy node + self.assertIsNotNone(req_ex.url, "request executor URL did not initialize") + self.assertTrue(req_ex.url.startswith(valid_url), f"unexpected URL: {req_ex.url}") + + # And simple operations should succeed + with store2.open_session() as session: + session.store({"Name": "Jane"}, "users/2") + session.save_changes() + + def test_subscription_failover_with_invalid_dns(self): + invalid_host = "http://thisnamedoesnotexist:8080" + valid_url = self.store.urls[0] + + with DocumentStore(urls=[invalid_host, valid_url], database=self.store.database) as store2: + store2.conventions.disable_topology_updates = False + store2.initialize() + + # Create a subscription and ensure worker connects and receives items + sub_id = store2.subscriptions.create_for_class(User) + with store2.subscriptions.get_subscription_worker(SubscriptionWorkerOptions(sub_id), User) as worker: + got_item = Event() + + def _run(batch): + for item in batch.items: + if item.result is not None: + got_item.set() + + worker.run(_run) + + # Add a document so the subscription has something to send + with store2.open_session() as session: + session.store(User(name="SubUser")) + session.save_changes() + + self.assertTrue(got_item.wait(10))