From 8ae02a221171b7da4e86112c2a01630f022a46b0 Mon Sep 17 00:00:00 2001 From: digitie Date: Wed, 16 Sep 2026 07:37:13 +0900 Subject: [PATCH] fix: paginate getVilageFcst/getUltraSrtFcst/getUltraSrtNcst instead of failing on a second page _fetch_items requested one 1,000-row page and raised a non-retryable KmaParseError the moment the response reported a next page, rather than fetching it. getVilageFcst alone publishes up to a dozen categories over a 3-day, 3-hour-step forecast, and how many rows that produces depends on how many categories the issuing office actually populated for a given base_time -- comfortably under 1,000 most of the time, but not always. A downstream consumer (kor-travel-weather) observed this failing for several consecutive hourly runs at a time, clustered in the evening KST hours, going stale for hours before the next base_time happened to fit. _fetch_items now walks pageNo forward, accumulating items.item across pages, until has_next_page() says there is nothing left, with a _MAX_FETCH_PAGES=20 backstop in case an upstream response never reports itself as the last page (turning that into a clear, bounded error instead of an infinite request loop). The three call sites (now, forecast_short, _forecast_vilage) share this method already, so all three benefit uniformly; none of them exposed a page-size knob to fix this from the caller's side, so this had to be fixed here. No test previously exercised the has_next_page=True branch at all -- covered now with a two-page fetch and a page-cap test, both confirmed to fail against the pre-fix code with exactly the production error text ("KMA response has more items than the requested page size"). Co-Authored-By: Claude Sonnet 5 Claude-Session: https://claude.ai/code/session_01YChhZnHkrqPpfZ3NehtB58 --- CHANGELOG.md | 8 +++ src/kma/client.py | 115 +++++++++++++++++++++++++++++-------------- tests/test_client.py | 109 ++++++++++++++++++++++++++++++++++++++++ 3 files changed, 195 insertions(+), 37 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index e1d830f..1bd22c8 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -13,6 +13,14 @@ ### 수정 +- `getUltraSrtNcst`/`getUltraSrtFcst`/`getVilageFcst` 응답이 요청한 페이지 하나를 넘으면 + `KmaClient._fetch_items()`가 재시도 불가 `KmaParseError("KMA response has more items than + the requested page size")`로 즉시 실패하던 문제 수정. 단기예보(`getVilageFcst`)는 3일치를 + 3시간 간격 최대 십수 개 카테고리로 발표하므로 한 base_time에 실제로 발표된 카테고리 수에 따라 + 1,000행 페이지를 넘을 수 있고, 그때마다 정상 응답을 파싱하지 않고 실패시켰다(운영에서 저녁 시간대에 + 걸쳐 여러 시간 연속 실패가 관측됨). `_fetch_items()`가 이제 `pageNo`를 늘려가며 + `has_next_page()`가 거짓이 될 때까지 계속 가져와 합치고, 응답이 끝을 알리지 않는 경우를 대비해 + `_MAX_FETCH_PAGES=20`으로 상한을 둔다. - asyncio 전환 재검증을 위한 2인 적대적 리뷰어 서브에이전트(동시성/자원관리 관점, 보안/데이터 무결성 관점) 감사에서 발견·검증된 버그 수정: `ApiHubClient.aiter_pages()`/ `AsyncApiHubClient.iter_pages()`가 공용 `pagination.aiter_pages()` 헬퍼를 거치지 않고 diff --git a/src/kma/client.py b/src/kma/client.py index f25ca0b..8a71a5a 100644 --- a/src/kma/client.py +++ b/src/kma/client.py @@ -44,6 +44,12 @@ DEFAULT_BASE_URL = "https://apis.data.go.kr/1360000/VilageFcstInfoService_2.0" SERVICE_NAME = "VilageFcstInfoService_2.0" +#: Safety cap on how many pages one grid's forecast may take. At 1,000 rows +#: per page this has never needed more than one in practice; the cap exists to +#: turn a data.go.kr response that never reports itself as the last page into +#: a clear error instead of an unbounded loop. +_MAX_FETCH_PAGES = 20 + @dataclass(frozen=True) class _KmaBody: @@ -274,44 +280,53 @@ async def _fetch_items( nx: int, ny: int, ) -> _FetchedItems: - response = await self._request_with_metadata( - endpoint, - { - "base_date": base_date, - "base_time": base_time, - "nx": nx, - "ny": ny, - }, - ) - if has_next_page(response.body): - raise KmaParseError( - "KMA response has more items than the requested page size", - provider="data.go.kr", - endpoint=enum_value(endpoint), - failure_kind="parse", - retryable=False, - ) - try: - items = response.body["items"]["item"] - except (KeyError, TypeError) as exc: - raise KmaParseError( - "KMA response did not contain items.item", - provider="data.go.kr", - endpoint=enum_value(endpoint), - failure_kind="parse", - retryable=False, - ) from exc - if isinstance(items, Mapping): - return _FetchedItems([items], response.metadata) - if not isinstance(items, list): - raise KmaParseError( - "KMA response items.item was not a list", - provider="data.go.kr", - endpoint=enum_value(endpoint), - failure_kind="parse", - retryable=False, + """Fetch every page of one grid's response. + + `getVilageFcst` alone carries up to a dozen categories over a 3-day, + 3-hour-step forecast, and the item count that produces varies with how + many of those categories the office actually published for this base + time -- comfortably under our 1,000-row page most of the time, but not + always. Treating a second page as a hard, non-retryable parse error + (the previous behaviour here) turned an ordinary large response into a + failed run instead of one extra request, silently for however many + consecutive base times stayed over the line. + """ + endpoint_name = enum_value(endpoint) + items: list[Mapping[str, Any]] = [] + metadata: ResponseMetadata | None = None + page_no = 1 + while True: + response = await self._request_with_metadata( + endpoint, + { + "base_date": base_date, + "base_time": base_time, + "nx": nx, + "ny": ny, + "pageNo": page_no, + }, ) - return _FetchedItems(items, response.metadata) + if metadata is None: + # The paginated request differs from page to page only in + # pageNo; base_date/base_time/nx/ny -- everything this + # metadata actually records -- stay fixed, so the first page + # describes the whole fetch. + metadata = response.metadata + items.extend(_page_items(response.body, endpoint_name)) + if not has_next_page(response.body): + break + page_no += 1 + if page_no > _MAX_FETCH_PAGES: + raise KmaParseError( + f"KMA response did not finish paginating after " + f"{_MAX_FETCH_PAGES} pages", + provider="data.go.kr", + endpoint=endpoint_name, + failure_kind="parse", + retryable=False, + ) + assert metadata is not None + return _FetchedItems(items, metadata) async def _request( self, @@ -520,6 +535,32 @@ def _parse_kma_body(response: Any, endpoint_name: str, metadata: ResponseMetadat return _KmaBody(body, metadata) +def _page_items(body: Mapping[str, Any], endpoint_name: str) -> list[Mapping[str, Any]]: + """Extract one page's ``items.item`` list, raising the same parse errors + ``_fetch_items`` always has for a shape it doesn't recognise.""" + try: + items = body["items"]["item"] + except (KeyError, TypeError) as exc: + raise KmaParseError( + "KMA response did not contain items.item", + provider="data.go.kr", + endpoint=endpoint_name, + failure_kind="parse", + retryable=False, + ) from exc + if isinstance(items, Mapping): + return [items] + if not isinstance(items, list): + raise KmaParseError( + "KMA response items.item was not a list", + provider="data.go.kr", + endpoint=endpoint_name, + failure_kind="parse", + retryable=False, + ) + return items + + def _forecast_item( item: Mapping[str, Any], endpoint: str | KmaEndpoint, diff --git a/tests/test_client.py b/tests/test_client.py index a3cdb4a..2541050 100644 --- a/tests/test_client.py +++ b/tests/test_client.py @@ -116,6 +116,39 @@ def _payload(items: Any) -> dict[str, Any]: } +def _paged_payload( + items: Any, *, page_no: int, num_of_rows: int, total_count: int +) -> dict[str, Any]: + return { + "response": { + "header": {"resultCode": "00", "resultMsg": "NORMAL_SERVICE"}, + "body": { + "pageNo": page_no, + "numOfRows": num_of_rows, + "totalCount": total_count, + "items": {"item": items}, + }, + } + } + + +class PagedFakeSession: + """Answers with a different payload per ``pageNo``, keyed by the request.""" + + def __init__(self, payloads_by_page: dict[int, dict[str, Any]]) -> None: + self.payloads_by_page = payloads_by_page + self.calls: list[dict[str, Any]] = [] + + async def get(self, url: str, *, params: dict[str, Any], timeout: float) -> FakeResponse: + self.calls.append({"url": url, "params": params, "timeout": timeout}) + page_no = int(params["pageNo"]) + return FakeResponse(self.payloads_by_page[page_no]) + + @property + def requested_pages(self) -> list[int]: + return [int(call["params"]["pageNo"]) for call in self.calls] + + def _error_payload(code: str, message: str = "ERROR") -> dict[str, Any]: return { "response": { @@ -530,3 +563,79 @@ async def test_malformed_forecast_item_raises_parse_error() -> None: ) (await assert_raises(KmaParseError, lambda: client.forecast(nx=60, ny=127))) + + +async def test_a_second_page_is_fetched_rather_than_treated_as_a_parse_error() -> None: + """`getVilageFcst` alone carries a dozen categories over a 3-day, 3-hour + forecast; how many rows that produces varies with how much the office + published for that base time, and it does not always fit one page. + Treating a second page as an unrecoverable parse error turned an + ordinary large response into a failed run -- for however many + consecutive base times stayed over the line -- instead of one extra + request.""" + session = PagedFakeSession( + { + 1: _paged_payload( + [ + {"category": "T1H", "obsrValue": "18.4"}, + {"category": "REH", "obsrValue": "52"}, + ], + page_no=1, + num_of_rows=2, + total_count=3, + ), + 2: _paged_payload( + [{"category": "WSD", "obsrValue": "3.1"}], + page_no=2, + num_of_rows=2, + total_count=3, + ), + } + ) + client = KmaClient("decoded-key", session=session) + + snapshot = await client.now(nx=60, ny=127, when=datetime(2026, 4, 30, 14, 45, tzinfo=KST)) + + assert snapshot.temperature == 18.4 + assert snapshot.humidity == 52 + assert snapshot.wind_speed == 3.1 + assert session.requested_pages == [1, 2] + assert snapshot.metadata is not None + # The paginated request differs from page to page only in pageNo; the + # metadata describes what identifies the fetch (base_date/base_time/grid), + # which is the same value on every page. + assert snapshot.metadata.base_time == "1400" + + +async def test_pagination_gives_up_after_the_page_cap_rather_than_looping_forever() -> None: + """A response that never reports itself as the last page must fail + loudly and boundedly, not hang the run consuming pages one at a time.""" + + class NeverLastPageSession: + def __init__(self) -> None: + self.calls = 0 + + async def get( + self, url: str, *, params: dict[str, Any], timeout: float + ) -> FakeResponse: + self.calls += 1 + page_no = int(params["pageNo"]) + # totalCount always claims one more row than this page reports, + # so has_next_page is true no matter how many pages are fetched. + return FakeResponse( + _paged_payload( + [{"category": "T1H", "obsrValue": "18.4"}], + page_no=page_no, + num_of_rows=1, + total_count=page_no + 1, + ) + ) + + session = NeverLastPageSession() + client = KmaClient("decoded-key", session=session) + + exc = await assert_raises(KmaParseError, lambda: client.now(nx=60, ny=127)) + + assert "paginat" in str(exc).lower() + # Bounded: the client gave up rather than fetching pages indefinitely. + assert session.calls <= 21