Skip to content
Merged
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
8 changes: 8 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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()` 헬퍼를 거치지 않고
Expand Down
115 changes: 78 additions & 37 deletions src/kma/client.py
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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,
Expand Down
109 changes: 109 additions & 0 deletions tests/test_client.py
Original file line number Diff line number Diff line change
Expand Up @@ -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": {
Expand Down Expand Up @@ -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
Loading