diff --git a/src/openhound_github/helpers.py b/src/openhound_github/helpers.py index aed2b99..e6a4210 100644 --- a/src/openhound_github/helpers.py +++ b/src/openhound_github/helpers.py @@ -1,11 +1,13 @@ import logging import time -from typing import Optional +from collections.abc import Iterator +from typing import Any, Optional from urllib.parse import urlparse from dlt.common import jsonpath from dlt.sources.helpers import requests from dlt.sources.helpers.rest_client.auth import AuthConfigBase +from dlt.sources.helpers.rest_client.client import RESTClient from dlt.sources.helpers.rest_client.paginators import ( JSONResponseCursorPaginator, ) @@ -20,6 +22,22 @@ class GraphQLPaginationError(RuntimeError): pass +class AdaptiveGraphQLPageError(RuntimeError): + """A terminal GraphQL page failure with adaptive pagination context.""" + + def __init__( + self, + *, + cursor: str | None, + page_size: int, + error: BaseException, + ) -> None: + super().__init__(str(error)) + self.cursor = cursor + self.page_size = page_size + self.error = error + + def scim_skip_reason(exception: BaseException) -> str | None: """Return a user-facing reason for expected SCIM API unavailability.""" if not isinstance(exception, requests.HTTPError) or exception.response is None: @@ -56,6 +74,10 @@ def __init__( self.has_next_field = has_next_field self.allow_missing_page_info = allow_missing_page_info + @property + def next_cursor(self) -> str | None: + return self._next_reference + def init_request(self, request: "Request") -> None: self._next_reference = None self._has_next_page = True @@ -149,6 +171,156 @@ def update_request(self, request: "Request") -> None: variables[self.cursor_variable] = self._next_reference +def _graphql_gateway_status(exception: BaseException) -> int | None: + if not isinstance(exception, requests.HTTPError) or exception.response is None: + return None + + status_code = exception.response.status_code + if status_code in (502, 504): + return status_code + return None + + +def adaptive_graphql_paginate( + client: RESTClient, + *, + query: str, + variables: dict[str, Any], + page_info_path: str, + cursor_variable: str = "after", + page_size_variable: str = "count", + page_sizes: tuple[int, ...] = (100, 50, 25), + allow_missing_page_info: bool = False, + resource_name: str = "graphql", + scope_name: str | None = None, + scope_value: str | None = None, + log: logging.Logger | None = None, +) -> Iterator[dict[str, Any]]: + """Yield GraphQL pages while adapting connection size after gateway failures. + + The underlying HTTP session still owns ordinary request retries. This helper + handles the remaining terminal 502/504 case by retrying the same cursor with + progressively smaller connection sizes. + """ + if not page_sizes or any(page_size <= 0 for page_size in page_sizes): + raise ValueError("page_sizes must contain positive integers") + if tuple(sorted(set(page_sizes), reverse=True)) != page_sizes: + raise ValueError("page_sizes must be unique and in descending order") + + logger_instance = log or logger + base_page_size = page_sizes[0] + cursor = variables.get(cursor_variable) + if cursor is not None and not isinstance(cursor, str): + raise ValueError(f"{cursor_variable} must be a string or None") + + stable_page_size = base_page_size + probe_base_page_size = False + + while True: + if stable_page_size == base_page_size: + candidate_sizes = page_sizes + probing_base_page_size = False + elif probe_base_page_size: + candidate_sizes = (base_page_size,) + tuple( + page_size + for page_size in page_sizes + if page_size <= stable_page_size + ) + probing_base_page_size = True + else: + candidate_sizes = tuple( + page_size + for page_size in page_sizes + if page_size <= stable_page_size + ) + probing_base_page_size = False + + page_data: dict[str, Any] | None = None + paginator: GraphQLCursorPaginator | None = None + successful_page_size: int | None = None + + for index, page_size in enumerate(candidate_sizes): + request_variables = { + **variables, + cursor_variable: cursor, + page_size_variable: page_size, + } + try: + response = client.post( + "/graphql", + json={"query": query, "variables": request_variables}, + ) + response.raise_for_status() + response_json = response.json() + paginator = GraphQLCursorPaginator( + page_info_path=page_info_path, + cursor_variable=cursor_variable, + cursor_field="endCursor", + has_next_field="hasNextPage", + allow_missing_page_info=allow_missing_page_info, + ) + paginator.update_state(response) + page_data = response_json.get("data") + if not isinstance(page_data, dict): + raise GraphQLPaginationError( + "GraphQL response data must be an object" + ) + successful_page_size = page_size + break + except Exception as error: + status_code = _graphql_gateway_status(error) + next_page_size = ( + candidate_sizes[index + 1] + if status_code is not None and index + 1 < len(candidate_sizes) + else None + ) + if next_page_size is not None: + scope = ( + f" {scope_name} '{scope_value}'" + if scope_name and scope_value + else "" + ) + logger_instance.warning( + "Adaptive GraphQL pagination for resource '%s'%s retrying " + "cursor %r with page size %d after HTTP %d at page size %d", + resource_name, + scope, + cursor, + next_page_size, + status_code, + page_size, + ) + continue + raise AdaptiveGraphQLPageError( + cursor=cursor, + page_size=page_size, + error=error, + ) from error + + if page_data is None or paginator is None or successful_page_size is None: + raise GraphQLPaginationError("GraphQL page did not produce data") + + yield page_data + + if successful_page_size == base_page_size: + stable_page_size = base_page_size + probe_base_page_size = False + elif stable_page_size == base_page_size: + stable_page_size = successful_page_size + probe_base_page_size = True + elif probing_base_page_size: + stable_page_size = successful_page_size + probe_base_page_size = False + elif successful_page_size < stable_page_size: + stable_page_size = successful_page_size + probe_base_page_size = False + + if not paginator.has_next_page: + break + + cursor = paginator.next_cursor + + def _response_message(response: requests.Response) -> str: try: response_data = response.json() diff --git a/src/openhound_github/resources/organization.py b/src/openhound_github/resources/organization.py index 931eb86..0b7b1e7 100644 --- a/src/openhound_github/resources/organization.py +++ b/src/openhound_github/resources/organization.py @@ -22,7 +22,12 @@ TEAM_MEMBERS_OVERFLOW_QUERY, TEAMS_QUERY, ) -from openhound_github.helpers import GraphQLCursorPaginator, scim_skip_reason +from openhound_github.helpers import ( + AdaptiveGraphQLPageError, + GraphQLCursorPaginator, + adaptive_graphql_paginate, + scim_skip_reason, +) from openhound_github.main import app from openhound_github.models import ( ActionPermission, @@ -919,25 +924,17 @@ def repositories_graphql(ctx: SourceContext): repository_cursor: str | None = None emitted_repositories = 0 try: - paginator = GraphQLCursorPaginator( + for page_data in adaptive_graphql_paginate( + client, + query=REPO_REFS_QUERY, + variables={"login": org_name, "after": None}, page_info_path="data.organization.repositories.pageInfo", - cursor_variable="after", - cursor_field="endCursor", - has_next_field="hasNextPage", - ) - data = { - "query": REPO_REFS_QUERY, - "variables": {"login": org_name, "count": 100, "after": None}, - } - - for page_data in client.paginate( - "/graphql", - method="POST", - json=data, - paginator=paginator, - data_selector="data", + resource_name="repositories_graphql", + scope_name="organization", + scope_value=org_name, + log=logger, ): - repos_page = page_data[0]["organization"]["repositories"] + repos_page = page_data["organization"]["repositories"] for repo in repos_page["nodes"]: repo_record = {**repo} branch_rulesets = repo_record.pop("branchRulesets", None) or {} @@ -952,22 +949,29 @@ def repositories_graphql(ctx: SourceContext): if isinstance(page_info, dict): repository_cursor = page_info.get("endCursor") except Exception as e: + failure = e.error if isinstance(e, AdaptiveGraphQLPageError) else e + failure_cursor = ( + e.cursor if isinstance(e, AdaptiveGraphQLPageError) else repository_cursor + ) + page_size = e.page_size if isinstance(e, AdaptiveGraphQLPageError) else None logger.error( "Error in resource 'repositories_graphql' processing organization '%s' " - "at repository cursor %r after emitting %d repositories " + "at repository cursor %r with page size %r after emitting %d repositories " "(%s): %s", org_name, - repository_cursor, + failure_cursor, + page_size, emitted_repositories, - type(e).__name__, - e, + type(failure).__name__, + failure, extra={ "resource": "repositories_graphql", "phase": "resource_iteration", "org_name": org_name, - "repository_cursor": repository_cursor, + "repository_cursor": failure_cursor, + "page_size": page_size, "emitted_repositories": emitted_repositories, - "error_type": type(e).__name__, + "error_type": type(failure).__name__, }, ) continue diff --git a/tests/test_repository_rulesets.py b/tests/test_repository_rulesets.py index 8ffdff5..d7be6d7 100644 --- a/tests/test_repository_rulesets.py +++ b/tests/test_repository_rulesets.py @@ -1,7 +1,10 @@ import duckdb +import json import logging from unittest.mock import MagicMock +import requests + from openhound_github.lookup import GithubLookup from openhound_github.models.repository import Repository from openhound_github.resources.organization import ( @@ -12,18 +15,20 @@ class _FakeClient: - def __init__(self) -> None: - self.request_cursors: list[str | None] = [] - - def paginate(self, *args, **kwargs): - self.request_cursors.append(kwargs["json"]["variables"]["after"]) - return iter( - [ - [ - _repository_page_data("R_1", "repo", branch_ruleset_count=2) - ] - ] - ) + def __init__(self, *responses: requests.Response | BaseException) -> None: + self.responses = list(responses) or [ + _graphql_response(_repository_page_data("R_1", "repo", branch_ruleset_count=2)) + ] + self.request_variables: list[dict[str, object]] = [] + + def post(self, _path: str, *, json: dict[str, object]): + variables = json["variables"] + assert isinstance(variables, dict) + self.request_variables.append({**variables}) + response = self.responses.pop(0) + if isinstance(response, BaseException): + raise response + return response def _repository_page_data( @@ -60,49 +65,27 @@ def _repository_page_data( } -class _FailingSecondPageClient: - def __init__(self) -> None: - self.request_cursors: list[str | None] = [] - - def paginate(self, *args, **kwargs): - variables = kwargs["json"]["variables"] - self.request_cursors.append(variables["after"]) - yield [ - _repository_page_data( - "R_1", - "repo", - branch_ruleset_count=2, - repository_end_cursor="cursor-page-2", - repositories_has_next_page=True, - ) - ] - variables["after"] = "cursor-page-2" - self.request_cursors.append(variables["after"]) - raise ConnectionError("GraphQL page failed after retries") - - -class _TwoPageClient: - def __init__(self) -> None: - self.request_cursors: list[str | None] = [] - - def paginate(self, *args, **kwargs): - variables = kwargs["json"]["variables"] - pages = [ - _repository_page_data( - "R_1", - "repo-1", - branch_ruleset_count=2, - repository_end_cursor="cursor-page-2", - repositories_has_next_page=True, - ), - _repository_page_data("R_2", "repo-2", branch_ruleset_count=0), - ] - for page in pages: - self.request_cursors.append(variables["after"]) - yield [page] - page_info = page["organization"]["repositories"]["pageInfo"] - if page_info["hasNextPage"]: - variables["after"] = page_info["endCursor"] +def _graphql_response( + data: dict[str, object] | None = None, + *, + status_code: int = 200, + text: str | None = None, +) -> requests.Response: + response = requests.Response() + response.status_code = status_code + if text is not None: + response._content = text.encode("utf-8") + else: + response._content = json.dumps({"data": data or {}}).encode("utf-8") + response.request = requests.Request("POST", "https://api.github.com/graphql").prepare() + return response + + +def _request_pages(client: _FakeClient) -> list[tuple[object, object]]: + return [ + (variables["after"], variables["count"]) + for variables in client.request_variables + ] def _make_repository() -> Repository: @@ -163,7 +146,18 @@ def test_repositories_graphql_flattens_branch_ruleset_count() -> None: def test_repositories_graphql_logs_cursor_and_emitted_count_on_page_failure( caplog, ) -> None: - client = _FailingSecondPageClient() + client = _FakeClient( + _graphql_response( + _repository_page_data( + "R_1", + "repo", + branch_ruleset_count=2, + repository_end_cursor="cursor-page-2", + repositories_has_next_page=True, + ) + ), + ConnectionError("GraphQL page failed after retries"), + ) ctx = SourceContext( client=client, organizations=[OrgContext(client=client, org_name="org")], @@ -175,14 +169,26 @@ def test_repositories_graphql_logs_cursor_and_emitted_count_on_page_failure( assert len(rows) == 1 assert ( "Error in resource 'repositories_graphql' processing organization 'org' " - "at repository cursor 'cursor-page-2' after emitting 1 repositories " + "at repository cursor 'cursor-page-2' with page size 100 " + "after emitting 1 repositories " "(ConnectionError): GraphQL page failed after retries" ) in caplog.text - assert client.request_cursors == [None, "cursor-page-2"] + assert _request_pages(client) == [(None, 100), ("cursor-page-2", 100)] def test_repositories_graphql_emits_all_repository_pages() -> None: - client = _TwoPageClient() + client = _FakeClient( + _graphql_response( + _repository_page_data( + "R_1", + "repo-1", + branch_ruleset_count=2, + repository_end_cursor="cursor-page-2", + repositories_has_next_page=True, + ) + ), + _graphql_response(_repository_page_data("R_2", "repo-2", branch_ruleset_count=0)), + ) ctx = SourceContext( client=client, organizations=[OrgContext(client=client, org_name="org")], @@ -194,7 +200,166 @@ def test_repositories_graphql_emits_all_repository_pages() -> None: ("R_1", 2), ("R_2", 0), ] - assert client.request_cursors == [None, "cursor-page-2"] + assert _request_pages(client) == [(None, 100), ("cursor-page-2", 100)] + + +def test_repositories_graphql_retries_gateway_failure_with_smaller_page_size( + caplog, +) -> None: + client = _FakeClient( + _graphql_response( + _repository_page_data( + "R_1", + "repo-1", + repository_end_cursor="cursor-page-2", + repositories_has_next_page=True, + ) + ), + _graphql_response(status_code=502, text="bad gateway"), + _graphql_response( + _repository_page_data( + "R_2", + "repo-2", + repository_end_cursor="cursor-page-3", + repositories_has_next_page=True, + ) + ), + _graphql_response(_repository_page_data("R_3", "repo-3")), + ) + ctx = SourceContext( + client=client, + organizations=[OrgContext(client=client, org_name="org")], + ) + + with caplog.at_level(logging.WARNING, logger="openhound_github.resources.organization"): + rows = list(repositories_graphql.__wrapped__(ctx)) + + assert [row["id"] for row in rows] == ["R_1", "R_2", "R_3"] + assert _request_pages(client) == [ + (None, 100), + ("cursor-page-2", 100), + ("cursor-page-2", 50), + ("cursor-page-3", 100), + ] + assert ( + "retrying cursor 'cursor-page-2' with page size 50 after HTTP 502 " + "at page size 100" + ) in caplog.text + + +def test_repositories_graphql_retries_gateway_failure_down_to_25() -> None: + client = _FakeClient( + _graphql_response( + _repository_page_data( + "R_1", + "repo-1", + repository_end_cursor="cursor-page-2", + repositories_has_next_page=True, + ) + ), + _graphql_response(status_code=502, text="bad gateway"), + _graphql_response(status_code=504, text="gateway timeout"), + _graphql_response(_repository_page_data("R_2", "repo-2")), + ) + ctx = SourceContext( + client=client, + organizations=[OrgContext(client=client, org_name="org")], + ) + + rows = list(repositories_graphql.__wrapped__(ctx)) + + assert [row["id"] for row in rows] == ["R_1", "R_2"] + assert _request_pages(client) == [ + (None, 100), + ("cursor-page-2", 100), + ("cursor-page-2", 50), + ("cursor-page-2", 25), + ] + + +def test_repositories_graphql_logs_terminal_gateway_failure_at_smallest_page_size( + caplog, +) -> None: + client = _FakeClient( + _graphql_response( + _repository_page_data( + "R_1", + "repo-1", + repository_end_cursor="cursor-page-2", + repositories_has_next_page=True, + ) + ), + _graphql_response(status_code=502, text="bad gateway"), + _graphql_response(status_code=504, text="gateway timeout"), + _graphql_response(status_code=502, text="bad gateway"), + ) + ctx = SourceContext( + client=client, + organizations=[OrgContext(client=client, org_name="org")], + ) + + with caplog.at_level(logging.ERROR, logger="openhound_github.resources.organization"): + rows = list(repositories_graphql.__wrapped__(ctx)) + + assert [row["id"] for row in rows] == ["R_1"] + assert _request_pages(client) == [ + (None, 100), + ("cursor-page-2", 100), + ("cursor-page-2", 50), + ("cursor-page-2", 25), + ] + assert ( + "at repository cursor 'cursor-page-2' with page size 25 " + "after emitting 1 repositories" + ) in caplog.text + + +def test_repositories_graphql_stays_degraded_after_failed_probe() -> None: + client = _FakeClient( + _graphql_response( + _repository_page_data( + "R_1", + "repo-1", + repository_end_cursor="cursor-page-2", + repositories_has_next_page=True, + ) + ), + _graphql_response(status_code=502, text="bad gateway"), + _graphql_response( + _repository_page_data( + "R_2", + "repo-2", + repository_end_cursor="cursor-page-3", + repositories_has_next_page=True, + ) + ), + _graphql_response(status_code=502, text="bad gateway"), + _graphql_response( + _repository_page_data( + "R_3", + "repo-3", + repository_end_cursor="cursor-page-4", + repositories_has_next_page=True, + ) + ), + _graphql_response(_repository_page_data("R_4", "repo-4")), + ) + ctx = SourceContext( + client=client, + organizations=[OrgContext(client=client, org_name="org")], + ) + + rows = list(repositories_graphql.__wrapped__(ctx)) + + assert [row["id"] for row in rows] == ["R_1", "R_2", "R_3", "R_4"] + assert _request_pages(client) == [ + (None, 100), + ("cursor-page-2", 100), + ("cursor-page-2", 50), + ("cursor-page-3", 100), + ("cursor-page-3", 50), + ("cursor-page-4", 50), + ] def test_repository_node_surfaces_branch_ruleset_presence() -> None: