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
7 changes: 6 additions & 1 deletion .importlinter
Original file line number Diff line number Diff line change
Expand Up @@ -63,10 +63,15 @@ ignore_imports =
dataretrieval -> dataretrieval.ogc.chunking

[importlinter:contract:ogc-facade]
name = NGWMN consumes the OGC facade only, never its internals (ADR 0007)
name = Facade-only OGC consumers, never its internals (ADR 0007)
type = forbidden
; Listed per module rather than by package: most of ``waterdata`` (``ratings``,
; ``reference``, ``samples``, ``stats``, the package ``__init__``) legitimately
; imports ``ogc`` internals today, so only the modules that have earned the
; facade-only seam belong here.
source_modules =
dataretrieval.ngwmn
dataretrieval.waterdata.cql
; The wildcard is what makes this durable: a new ``ogc`` submodule is covered
; the day it is added, without editing this contract.
forbidden_modules =
Expand Down
14 changes: 7 additions & 7 deletions dataretrieval/ogc/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -4,22 +4,22 @@

- :class:`OgcDialect` — per-API request/response quirks.
- :func:`prepare_request_args` — normalize caller kwargs for the engine.
- :func:`get_ogc_data` — full orchestrated OGC fetch (chunking + pagination).
- :func:`fetch_ogc_request` — execute a pre-built request with pagination.
- :func:`get_ogc_data` — full orchestrated OGC fetch (chunking + pagination),
including verbatim-CQL2 queries via its ``cql_body`` parameter.

Collection adapters (NGWMN, Water Data's generic wrapper) import from this
facade rather than reaching into engine internals. Generic execution policy
lives in :mod:`dataretrieval.transport`; the engine retains compatibility
wrappers at previous private paths.
facade rather than reaching into engine internals — every name here is usable
through the facade alone. Generic execution policy lives in
:mod:`dataretrieval.transport`; the engine retains compatibility wrappers at
previous private paths.
"""

from dataretrieval.ogc.engine import fetch_ogc_request, get_ogc_data
from dataretrieval.ogc.engine import get_ogc_data
from dataretrieval.ogc.policy import OgcDialect
from dataretrieval.ogc.requests import prepare_request_args

__all__ = [
"OgcDialect",
"fetch_ogc_request",
"get_ogc_data",
"prepare_request_args",
]
30 changes: 0 additions & 30 deletions dataretrieval/ogc/context.py

This file was deleted.

137 changes: 78 additions & 59 deletions dataretrieval/ogc/engine.py
Original file line number Diff line number Diff line change
Expand Up @@ -36,7 +36,6 @@
import pandas as pd

import dataretrieval.ogc.chunking as chunking
from dataretrieval.ogc.context import _dialect, _ogc_base_url, _row_cap
from dataretrieval.ogc.errors import _raise_for_non_200
from dataretrieval.ogc.policy import (
DEFAULT_DIALECT,
Expand All @@ -48,6 +47,7 @@
# the symbols its orchestration uses.
from dataretrieval.ogc.requests import (
_construct_api_requests,
_construct_cql_request,
_switch_arg_id,
_switch_properties_id,
)
Expand Down Expand Up @@ -130,6 +130,7 @@ async def _paginate(
follow_up: Callable[[_Cursor, httpx.AsyncClient], Awaitable[httpx.Response]],
client: httpx.AsyncClient | None = None,
raise_for_status: Callable[[httpx.Response], None] = _raise_for_non_200,
row_cap: int | None = None,
) -> tuple[pd.DataFrame, httpx.Response]:
"""Compatibility wrapper around collection-neutral cursor pagination."""
session = client if client is not None else active_client()
Expand All @@ -139,7 +140,7 @@ async def _paginate(
follow_up=follow_up,
client=session,
raise_for_status=raise_for_status,
row_cap=_row_cap.get(),
row_cap=row_cap,
)


Expand All @@ -164,6 +165,8 @@ async def _walk_pages(
geopd: bool,
req: httpx.Request,
client: httpx.AsyncClient | None = None,
*,
row_cap: int | None = None,
) -> tuple[pd.DataFrame, httpx.Response]:
"""
Iterate paginated OGC API responses and aggregate them into one DataFrame.
Expand All @@ -182,6 +185,11 @@ async def _walk_pages(
client : httpx.AsyncClient, optional
Caller-borrowed client; ``None`` defers client management to
:func:`_paginate`.
row_cap : int, optional
Stop following pages once this many rows have accumulated and
truncate to exactly this many. ``None`` (default) walks every page.
An early-stop download bound only — the combined-result cap is
applied in :func:`~dataretrieval.ogc.shaping._finalize_ogc`.

Returns
-------
Expand Down Expand Up @@ -211,6 +219,7 @@ async def follow_up(cursor: str, sess: httpx.AsyncClient) -> httpx.Response:
parse_response=functools.partial(_ogc_parse_response, geopd=geopd),
follow_up=follow_up,
client=client,
row_cap=row_cap,
)


Expand All @@ -223,6 +232,7 @@ def get_ogc_data(
max_rows: int | None = None,
extra_id_cols: frozenset[str] | set[str] = frozenset(),
dialect: OgcDialect | None = None,
cql_body: str | None = None,
) -> tuple[pd.DataFrame, BaseMetadata]:
"""
Retrieves OGC (Open Geospatial Consortium) data as a DataFrame with metadata.
Expand Down Expand Up @@ -262,6 +272,13 @@ def get_ogc_data(
dialect : OgcDialect, optional
Per-API request quirks (CQL2-only collections, date-only collections).
Defaults to a plain OGC API with neither.
cql_body : str, optional
A verbatim CQL2 JSON body to POST against ``collection`` instead of
building the query from ``args``. With a body, only the
``properties``, ``bbox``, ``limit``, ``skip_geometry``, and
``convert_type`` keys of ``args`` are consulted; there are no
multi-value axes to chunk, and the body's size is the server's
judgement, mirroring the planner's cql-json passthrough.

Returns
-------
Expand Down Expand Up @@ -303,7 +320,8 @@ def get_ogc_data(
# this function). ``_finalize_ogc`` is the single source of result shape;
# it also applies ``max_rows`` to the *combined* frame so the cap is the
# exact total even when the plan chunks or the call is resumed, while
# ``_row_cap`` below only early-stops each chunk's pagination.
# the per-chunk ``row_cap`` bound below only early-stops each chunk's
# pagination.
finalize = functools.partial(
_finalize_ogc,
properties=properties,
Expand All @@ -315,70 +333,71 @@ def get_ogc_data(
dialect=dialect,
base_url=base_url,
)

if cql_body is not None:
# ``args["properties"]`` holds the wire property list after the
# id-switch above; ``finalize`` holds the pre-switch, user-facing
# list, exactly as on the chunked path.
req = _construct_cql_request(
collection,
cql_body,
base_url=base_url,
properties=args.get("properties"),
bbox=args.get("bbox"),
limit=args.get("limit"),
skip_geometry=args.get("skip_geometry"),
)

# A one-item fan-out: same executor, retry, progress line, and
# resumable-interruption semantics as the chunked path — and the same
# ``finalize``, so a later ``exc.call.resume()`` returns the finished
# ``(df, BaseMetadata)`` shape rather than a raw response pair.
return FanOut(
[req],
functools.partial(_walk_pages, GEOPANDAS, row_cap=max_rows),
RetryPolicy.from_env(),
finalize,
canonical_url=str(req.url),
service=collection,
).resume()

# Bind the API target and quirks into the request builder and fetcher the
# same way ``finalize`` binds its own state: with ``functools.partial``.
# The plan sizes candidate chunks and a later ``exc.call.resume()``
# rebuilds them through these same bound callables, so the values the
# call was created with reach every chunk — even a resume fired long
# after this function returned — without any ambient state to snapshot.
build_request = functools.partial(
_construct_api_requests, base_url=base_url, dialect=dialect
)
fetch = functools.partial(
_fetch_once, build_request=build_request, row_cap=max_rows
)
run = chunking.multi_value_chunked(build_request=build_request)(fetch)
# No progress block here: the executor that emits the events owns the line
# (see :meth:`~dataretrieval.transport.fanout.FanOut.resume`).
with _row_cap(max_rows), _ogc_base_url(base_url), _dialect(dialect):
return _fetch_once(args, finalize=finalize)
return run(args, finalize=finalize)


@chunking.multi_value_chunked(build_request=_construct_api_requests)
async def _fetch_once(
args: dict[str, Any],
*,
build_request: Callable[..., httpx.Request],
row_cap: int | None = None,
) -> tuple[pd.DataFrame, httpx.Response]:
"""Send one prepared-args OGC request asynchronously; return (frame, response).

``@chunking.multi_value_chunked`` models every multi-value list
parameter and the cql-text filter as a chunkable axis, greedy-halves
the biggest chunk across all axes until each chunk URL fits,
and iterates the cartesian product. With no chunkable inputs the
decorator passes args through unchanged. The decorator gathers every
The undecorated per-chunk fetcher: ``get_ogc_data`` binds
``build_request`` (the target API's request builder) and ``row_cap``,
then wraps the result in ``chunking.multi_value_chunked``, which models
every multi-value list parameter and the cql-text filter as a chunkable
axis, greedy-halves the biggest chunk across all axes until each chunk
URL fits, and iterates the cartesian product. With no chunkable inputs
the decorator passes args through unchanged. The decorator gathers every
chunk over one shared :class:`httpx.AsyncClient` (concurrency
bounded by a semaphore, sized from ``API_USGS_CONCURRENT``). It also
returns a *synchronous* wrapper, so ``get_ogc_data`` keeps calling
``_fetch_once(args, finalize=...)`` synchronously. The return shape is
``(frame, response)``.
bounded by a semaphore, sized from ``API_USGS_CONCURRENT``) and
returns a *synchronous* wrapper, so ``get_ogc_data`` drives it
synchronously. The return shape is ``(frame, response)``.
"""
req = _construct_api_requests(**args)
return await _walk_pages(geopd=GEOPANDAS, req=req)


def fetch_ogc_request(
request: httpx.Request,
*,
collection: str,
) -> tuple[pd.DataFrame, httpx.Response]:
"""Execute a prepared OGC request with pagination, returning (df, response).

This is the facade-level entry point for generalized CQL requests: the
caller builds its own :class:`httpx.Request` (e.g. via
:func:`~dataretrieval.ogc.requests._construct_cql_request`) and hands it
here. The request is driven as a one-item
:class:`~dataretrieval.transport.fanout.FanOut` -- the same executor the
typed getters use -- so pagination, retry, progress reporting, and error
handling are identical to their path through :func:`_walk_pages`.

Parameters
----------
request : httpx.Request
A fully-constructed OGC API request (typically a POST/CQL2).
collection : str
Collection name, used only for progress-context labelling.

Returns
-------
pd.DataFrame
Concatenated page results.
httpx.Response
Aggregated response metadata.
"""

async def _fetch(req: httpx.Request) -> tuple[pd.DataFrame, httpx.Response]:
return await _walk_pages(geopd=GEOPANDAS, req=req)

return FanOut(
[request],
_fetch,
RetryPolicy.from_env(),
canonical_url=str(request.url),
service=collection,
).resume()
req = build_request(**args)
return await _walk_pages(geopd=GEOPANDAS, req=req, row_cap=row_cap)
5 changes: 2 additions & 3 deletions dataretrieval/ogc/policy.py
Original file line number Diff line number Diff line change
Expand Up @@ -5,9 +5,8 @@
It depends only on the stdlib, so any OGC submodule can import it without
creating cycles.

It names no endpoint: which collection an OGC call targets is the *adapter's*
policy, supplied per call as ``base_url`` (see
:data:`dataretrieval.ogc.context._ogc_base_url`). A default here would quietly
It names no endpoint: which API an OGC call targets is the *adapter's*
policy, supplied per call as ``base_url``. A default here would quietly
point every generic OGC caller at one API.

It must NOT import engine, shaping, or any collection adapter.
Expand Down
44 changes: 29 additions & 15 deletions dataretrieval/ogc/requests.py
Original file line number Diff line number Diff line change
@@ -1,10 +1,11 @@
"""OGC argument normalization and HTTP request construction.

Ambient request state lives in :mod:`dataretrieval.ogc.context`; queryables and
schema execution live in :mod:`dataretrieval.ogc.schema`. Neither is re-exported
from here -- importing the schema helper only to forward it would give this
module an edge to the one part of OGC that executes HTTP, which is exactly what
request *construction* is supposed to be free of.
The API to target and its quirks are explicit parameters (``base_url``,
``dialect``) -- construction states everything it needs. Queryables and schema
execution live in :mod:`dataretrieval.ogc.schema`, not re-exported from here --
importing the schema helper only to forward it would give this module an edge
to the one part of OGC that executes HTTP, which is exactly what request
*construction* is supposed to be free of.
"""

from __future__ import annotations
Expand All @@ -16,9 +17,8 @@

import httpx

from dataretrieval.ogc.context import _dialect as _dialect
from dataretrieval.ogc.context import _ogc_base_url as _ogc_base_url
from dataretrieval.ogc.dates import _DATE_RANGE_PARAMS, _format_api_dates
from dataretrieval.ogc.policy import DEFAULT_DIALECT, OgcDialect
from dataretrieval.transport.http import default_headers as _default_headers

# ---------------------------------------------------------------------------
Expand Down Expand Up @@ -120,9 +120,9 @@ def _partition_request_params(
return get_params, {}


def _items_url(collection: str) -> str:
"""The OGC items endpoint for ``collection`` under the active base URL."""
return f"{_ogc_base_url.get()}/collections/{collection}/items"
def _items_url(collection: str, base_url: str) -> str:
"""The OGC items endpoint for ``collection`` under ``base_url``."""
return f"{base_url}/collections/{collection}/items"


def _cql2_post_request(
Expand All @@ -146,11 +146,20 @@ def _construct_api_requests(
bbox: list[float] | None = None,
limit: int | None = None,
skip_geometry: bool | None = None,
*,
base_url: str,
dialect: OgcDialect | None = None,
**kwargs: Any,
) -> httpx.Request:
"""Construct an HTTP request object for the specified OGC API collection."""
service_url = _items_url(collection)
dialect = _dialect.get()
"""Construct an HTTP request object for the specified OGC API collection.

``base_url`` is required: this package is API-neutral and names no API of
its own, so the adapter naming the collection states the API it targets.
``dialect`` defaults to a plain OGC API with no per-collection quirks.
"""
service_url = _items_url(collection, base_url)
if dialect is None:
dialect = DEFAULT_DIALECT
for key in _DATE_RANGE_PARAMS:
if key in kwargs:
kwargs[key] = _format_api_dates(
Expand Down Expand Up @@ -190,13 +199,18 @@ def _construct_cql_request(
collection: str,
cql_body: str,
*,
base_url: str,
properties: list[str] | None = None,
bbox: list[float] | None = None,
limit: int | None = None,
skip_geometry: bool | None = None,
) -> httpx.Request:
"""Build a POST/CQL2 request from a verbatim CQL2 body."""
service_url = _items_url(collection)
"""Build a POST/CQL2 request from a verbatim CQL2 body.

``base_url`` is required for the same reason as in
:func:`_construct_api_requests`.
"""
service_url = _items_url(collection, base_url)
params = _ogc_query_params(
{},
properties=properties,
Expand Down
Loading