Skip to content

aeolus.schema, aeolus.qa, aeolus.network_registry

The v0.5.0 contract: the public column set, the QA enums, and the registry that maps each network's own quality tokens onto them.

Schema

The public wire format, and the one step that produces it.

Adapters emit types.ADAPTER_DATA_COLUMNS. finalise_data_frame runs where every download converges and adds the network identity and the QA model. It is the only code that knows both schemas — do not add these columns in adapters.

DATA_COLUMNS = ['site_code', 'network', 'date_time', 'measurand', 'value', 'units', 'qa_code', 'qa_tier', 'ratification_stage', 'backend', 'source_network', 'ratification', 'created_at'] module-attribute

METADATA_COLUMNS = ['site_code', 'site_name', 'latitude', 'longitude', 'network', 'country', 'instrument_class', 'provider', 'backend', 'measurands', 'source_network'] module-attribute

public_data_columns()

The column list in force: with or without the legacy mirrors.

Source code in src/aeolus/schema.py
def public_data_columns() -> list[str]:
    """The column list in force: with or without the legacy mirrors."""
    if options.legacy_columns:
        return list(DATA_COLUMNS)
    return [c for c in DATA_COLUMNS if c not in LEGACY_DATA_COLUMNS]

finalise_data_frame(df, source)

Turn an adapter frame for source into the public schema. Idempotent.

Source code in src/aeolus/schema.py
def finalise_data_frame(df: pd.DataFrame, source: str) -> pd.DataFrame:
    """Turn an adapter frame for *source* into the public schema. Idempotent."""
    attrs = dict(df.attrs)
    if df.empty:
        out = empty_public_frame()
        out.attrs = attrs
        return out

    network, backend, spec = spec_for_source(source)
    out = df.copy()
    out["network"] = network
    # An adapter that serves one network from several upstream datasets (EEA)
    # says which served each row, as a refinement of the route's backend
    # (EEA -> EEA_E1A). Anything else is the route's backend.
    if "backend" in out.columns:
        given = out["backend"]
        name = given.astype("string")
        refines = (name.notna() & ((name == backend) | name.str.startswith(f"{backend}_"))).fillna(False).astype(bool)
        out["backend"] = given.where(refines, backend)  # where() keeps the dtype: finalising twice must be a no-op
    else:
        out["backend"] = backend
    out["source_network"] = network

    wired = "qa_code" in out.columns
    codes = out["qa_code"] if wired else pd.Series(None, index=out.index, dtype=object)
    out["qa_code"] = codes.astype(object).where(codes.notna(), None)
    out["qa_tier"], out["ratification_stage"] = derive_qa(
        out["qa_code"], spec.qa_code_vocabulary, spec.default_ratification_stage
    )
    if wired or "ratification" not in out.columns:
        out["ratification"] = legacy_ratification(out["qa_tier"], out["ratification_stage"])

    global _warned_legacy
    if options.legacy_columns and not _warned_legacy:
        _warned_legacy = True
        warnings.warn(
            "The `source_network` and `ratification` columns are deprecated mirrors of "
            "`network` and of `qa_code`/`qa_tier`/`ratification_stage`; they will be removed "
            "in aeolus 1.0. Set AEOLUS_LEGACY_COLUMNS=0 (or aeolus.options.legacy_columns = "
            "False) to drop them now and check your code has migrated.",
            DeprecationWarning,
            stacklevel=4,
        )
    out = out[public_data_columns()]
    out.attrs = attrs
    return out

finalise_metadata_frame(df, source)

Add network identity to an adapter's site metadata. Extra columns are kept.

Source code in src/aeolus/schema.py
def finalise_metadata_frame(df: pd.DataFrame, source: str) -> pd.DataFrame:
    """Add network identity to an adapter's site metadata. Extra columns are kept."""
    network, backend, spec = spec_for_source(source)
    out = df.copy()
    for column in ("site_name", "latitude", "longitude", "measurands"):
        if column not in out.columns:
            out[column] = None
    out["network"] = network
    out["source_network"] = network
    out["country"] = spec.country
    out["backend"] = backend
    if "instrument_class" not in out.columns:
        out["instrument_class"] = spec.instrument_class
    # LMAM's provider is its pcode subfolder (spec D4); null for every other network
    out["provider"] = out["pcode"] if network == "LMAM" and "pcode" in out.columns else None
    core = public_metadata_columns()
    extras = [c for c in out.columns if c not in METADATA_COLUMNS]
    return out[core + extras]

QA enums and derivation

The v0.5.0 data-quality model: frozen enums and pure derivation functions.

The enum values are permanent public keys (renaming one is a major breaking change), shared with Argus. qa_code is whatever the upstream wrote; qa_tier and ratification_stage are derived from it through the network's vocabulary (see aeolus.network_registry).

QA_TIERS = ('reference_full_qc', 'reference_provisional', 'lcs_calibrated', 'lcs_factory_only', 'flagged', 'unknown') module-attribute

RATIFICATION_STAGES = ('unratified', 'ratified', 'supplied', 'not_applicable') module-attribute

QA_MODELS = ('regulatory_temporal_ratification', 'staged_calibration', 'point_in_time_validation', 'none', 'mixed', 'unknown') module-attribute

INSTRUMENT_CLASSES = ('reference', 'equivalent', 'indicative', 'LCS', 'mixed', 'unknown') module-attribute

derive_qa(qa_codes, vocabulary, default_stage)

Return (qa_tier, ratification_stage) for a series of upstream codes.

A missing code means the upstream said nothing: tier unknown and the network's default_stage. A code the vocabulary does not list is also unknown, with a null stage — an unrecognised token must not be silently promoted to the network default.

Source code in src/aeolus/qa.py
def derive_qa(
    qa_codes: pd.Series, vocabulary: dict[str, dict], default_stage: str | None
) -> tuple[pd.Series, pd.Series]:
    """Return ``(qa_tier, ratification_stage)`` for a series of upstream codes.

    A missing code means the upstream said nothing: tier ``unknown`` and the
    network's *default_stage*. A code the vocabulary does not list is also
    ``unknown``, with a null stage — an unrecognised token must not be
    silently promoted to the network default.
    """
    codes = qa_codes.astype(object)
    missing = codes.isna()
    tier_map = {code: entry["qa_tier"] for code, entry in vocabulary.items()}
    stage_map = {code: entry.get("ratification_stage") for code, entry in vocabulary.items()}

    tier = codes.map(tier_map).where(~missing, "unknown").fillna("unknown")
    stage = codes.map(stage_map).astype(object)
    stage = stage.where(~missing, default_stage)
    return tier.astype(object), stage.where(stage.notna(), None)

legacy_ratification(qa_tier, stage)

The pre-0.5.0 ratification string, derived from the new columns (spec §4.3).

Source code in src/aeolus/qa.py
def legacy_ratification(qa_tier: pd.Series, stage: pd.Series) -> pd.Series:
    """The pre-0.5.0 ``ratification`` string, derived from the new columns (spec §4.3)."""
    out = pd.Series("None", index=qa_tier.index, dtype=object)
    out = out.where(~(qa_tier == "unknown") | stage.isna(), "Unvalidated")
    for tier, label in _LEGACY_BY_TIER.items():
        out = out.where(qa_tier != tier, label)
    out = out.where(stage != "ratified", "Ratified")
    out = out.where(stage != "unratified", "Provisional")
    return out

Network registry

Each network's identity, QA vocabulary, when its data gets ratified, and its licence live in src/aeolus/data/qa_vocabularies/<CODE>.yaml and are read through this module.

Network registry: who each network is and how to read its QA codes.

A network is who produced the data (AURN, LAQN, ...). A source is a way aeolus fetches it (AURN via RData, AURN-SOS via the SOS API). Several sources can serve one network; SOURCE_ROUTES records which. The routing engine (choosing a backend) is v0.6.0 — this is only the lookup table.

SOURCE_ROUTES = {**{n: (n, 'RDATA') for n in ('AURN', 'SAQN', 'WAQN', 'NI', 'AQE', 'LMAM', 'LAQN')}, 'SAQD': ('SAQN', 'RDATA'), **{f'{n}-SOS': (n, 'SOS') for n in ('AURN', 'SAQN', 'WAQN', 'NI', 'AQE')}, 'LAQN-ERG': ('LAQN', 'ERG_REST'), **{n: (n, n) for n in ('BREATHE_LONDON', 'AIRQO', 'AIRNOW', 'EEA', 'SONITUS', 'SENSOR_COMMUNITY', 'OPENAQ', 'PURPLEAIR')}} module-attribute

NetworkSpec dataclass

Source code in src/aeolus/network_registry.py
@dataclass(frozen=True)
class NetworkSpec:
    code: str
    name: str
    country: str
    regulatory: bool
    instrument_class: str
    qa_model: str
    default_ratification_stage: str | None
    qa_code_vocabulary: dict = field(default_factory=dict)
    operators: tuple = ()
    ratification_cadence_months: int | None = None
    first_ratification_latency_months: int | None = None
    full_year_ratified_by: str | None = None
    ratification_overrides: str = ""
    data_licence: str = ""
    homepage_url: str = ""
    notes: str = ""

code instance-attribute

name instance-attribute

country instance-attribute

regulatory instance-attribute

instrument_class instance-attribute

qa_model instance-attribute

default_ratification_stage instance-attribute

qa_code_vocabulary = field(default_factory=dict) class-attribute instance-attribute

operators = () class-attribute instance-attribute

ratification_cadence_months = None class-attribute instance-attribute

first_ratification_latency_months = None class-attribute instance-attribute

full_year_ratified_by = None class-attribute instance-attribute

ratification_overrides = '' class-attribute instance-attribute

data_licence = '' class-attribute instance-attribute

homepage_url = '' class-attribute instance-attribute

notes = '' class-attribute instance-attribute

__init__(code, name, country, regulatory, instrument_class, qa_model, default_ratification_stage, qa_code_vocabulary=dict(), operators=(), ratification_cadence_months=None, first_ratification_latency_months=None, full_year_ratified_by=None, ratification_overrides='', data_licence='', homepage_url='', notes='')

get_network_spec(code)

Source code in src/aeolus/network_registry.py
def get_network_spec(code: str) -> NetworkSpec:
    try:
        return _load()[code.upper()]
    except KeyError:
        raise KeyError(f"Unknown network {code.upper()!r}. Known: {sorted(_load())}") from None

list_network_specs()

Source code in src/aeolus/network_registry.py
def list_network_specs() -> list[NetworkSpec]:
    return list(_load().values())

route_for(source)

Source code in src/aeolus/network_registry.py
def route_for(source: str) -> tuple[str, str]:
    try:
        return SOURCE_ROUTES[source.upper()]
    except KeyError:
        raise KeyError(
            f"Source {source.upper()!r} has no entry in network_registry.SOURCE_ROUTES"
        ) from None

spec_for_source(source)

(network, backend, spec) for a source — including one aeolus has never heard of.

A source registered with register_source but absent from SOURCE_ROUTES (a user's own adapter) is treated as its own network with nothing known about it: qa_tier is unknown and no stage is assumed.

Source code in src/aeolus/network_registry.py
def spec_for_source(source: str) -> tuple[str, str, NetworkSpec]:
    """``(network, backend, spec)`` for a source — including one aeolus has never
    heard of.

    A source registered with ``register_source`` but absent from
    ``SOURCE_ROUTES`` (a user's own adapter) is treated as its own network with
    nothing known about it: ``qa_tier`` is ``unknown`` and no stage is assumed.
    """
    name = source.upper()
    if name in SOURCE_ROUTES:
        network, backend = SOURCE_ROUTES[name]
        return network, backend, get_network_spec(network)
    return name, name, NetworkSpec(
        code=name, name=name, country="*", regulatory=False,
        instrument_class="unknown", qa_model="unknown", default_ratification_stage=None,
    )

Units and cache

One spelling per unit.

Upstreams write the same unit many ways (ug.m-3, µg/m³, UG/M3). Every adapter and every metrics path canonicalises through :func:canonical_unit, so there is a single list to maintain rather than one per module.

canonical_unit(unit)

Return the canonical spelling of unit (ug/m3, mg/m3, ng/m3, ppb, ppm).

Units that are not concentration units (C, %, hPa, AQI) and missing values are returned unchanged.

Source code in src/aeolus/units.py
def canonical_unit(unit):
    """Return the canonical spelling of *unit* (``ug/m3``, ``mg/m3``, ``ng/m3``, ``ppb``, ``ppm``).

    Units that are not concentration units (``C``, ``%``, ``hPa``, ``AQI``) and
    missing values are returned unchanged.
    """
    if unit is None or not isinstance(unit, str):
        return unit
    text = unit.strip()
    folded = text.lower().replace("µ", "u").replace("μ", "u").replace("³", "3")
    match = _MASS_PER_VOLUME.match(folded)
    if match:
        return f"{match.group(1)}g/m3"
    if folded in ("ppb", "ppm"):
        return folded
    return text

canonical_units(units)

Vectorised :func:canonical_unit, mapping each distinct value once.

Source code in src/aeolus/units.py
def canonical_units(units: pd.Series) -> pd.Series:
    """Vectorised :func:`canonical_unit`, mapping each distinct value once."""
    distinct = pd.unique(units.astype(object))
    return units.astype(object).map({u: canonical_unit(u) for u in distinct})

Local file cache for downloaded air quality data.

Caches data as Parquet files, keyed by source, site, and date range. This avoids redundant API calls when re-running notebooks or analyses.

Cache location defaults to ~/.cache/aeolus/ and can be overridden by setting the AEOLUS_CACHE_DIR environment variable.

Complete results for an explicit date range never expire. Two kinds of entry are volatile and are re-fetched once older than AEOLUS_CACHE_VOLATILE_TTL_S seconds (default 3600): rolling last= windows, which are keyed on the shorthand so a re-run hits the cache, and results missing a requested site, which may reflect a transient failure. A last= window no longer than the TTL is always fetched live.

Usage::

import aeolus
from aeolus.cache import enable_cache, disable_cache, clear_cache

# Enable caching (all subsequent downloads are cached)
enable_cache()

# Downloads hit the API on first call, then use cache
data = aeolus.download("AURN", ["MY1"], start, end)
data = aeolus.download("AURN", ["MY1"], start, end)  # instant

# Clear everything
clear_cache()

# Disable caching
disable_cache()

enable_cache(cache_dir=None)

Enable local file caching for downloads.

Subsequent calls to aeolus.download() will check the cache before hitting the network. Cached data is stored as Parquet files.

Parameters:

Name Type Description Default
cache_dir str | Path | None

Override the cache directory. Defaults to ~/.cache/aeolus/ or AEOLUS_CACHE_DIR env var.

None

Example::

>>> import aeolus
>>> from aeolus.cache import enable_cache
>>> enable_cache()
>>> data = aeolus.download("AURN", ["MY1"], start, end)  # fetches
>>> data = aeolus.download("AURN", ["MY1"], start, end)  # cached
Source code in src/aeolus/cache.py
def enable_cache(cache_dir: str | Path | None = None) -> None:
    """
    Enable local file caching for downloads.

    Subsequent calls to ``aeolus.download()`` will check the cache before
    hitting the network. Cached data is stored as Parquet files.

    Args:
        cache_dir: Override the cache directory. Defaults to
                   ``~/.cache/aeolus/`` or ``AEOLUS_CACHE_DIR`` env var.

    Example::

        >>> import aeolus
        >>> from aeolus.cache import enable_cache
        >>> enable_cache()
        >>> data = aeolus.download("AURN", ["MY1"], start, end)  # fetches
        >>> data = aeolus.download("AURN", ["MY1"], start, end)  # cached
    """
    global _cache_enabled, _cache_dir
    _cache_enabled = True
    if cache_dir is not None:
        _cache_dir = Path(cache_dir)
    logger.info("Cache enabled: %s", _get_cache_dir())

disable_cache()

Disable local file caching.

Downloads will always go to the network. Existing cache files are preserved (use clear_cache() to remove them).

Source code in src/aeolus/cache.py
def disable_cache() -> None:
    """
    Disable local file caching.

    Downloads will always go to the network. Existing cache files
    are preserved (use ``clear_cache()`` to remove them).
    """
    global _cache_enabled
    _cache_enabled = False
    logger.info("Cache disabled")

clear_cache(source=None)

Remove cached files.

Parameters:

Name Type Description Default
source str | None

If given, only clear cache for this source. Otherwise clears the entire cache.

None

Returns:

Type Description
int

Number of files removed.

Example::

>>> from aeolus.cache import clear_cache
>>> clear_cache("AURN")       # clear AURN cache only
>>> clear_cache()             # clear everything
Source code in src/aeolus/cache.py
def clear_cache(source: str | None = None) -> int:
    """
    Remove cached files.

    Args:
        source: If given, only clear cache for this source.
                Otherwise clears the entire cache.

    Returns:
        Number of files removed.

    Example::

        >>> from aeolus.cache import clear_cache
        >>> clear_cache("AURN")       # clear AURN cache only
        >>> clear_cache()             # clear everything
    """
    cache_dir = _get_cache_dir()
    count = 0

    if source:
        # Current entries, plus any written by an earlier cache version
        for target in (_entries_dir() / source.upper(), cache_dir / source.upper()):
            if target.exists():
                for f in target.glob("*.parquet"):
                    f.unlink(missing_ok=True)
                    count += 1
                # Remove empty directory
                if not any(target.iterdir()):
                    target.rmdir()
    else:
        for f in cache_dir.rglob("*.parquet"):
            f.unlink(missing_ok=True)
            count += 1
        # Remove empty subdirectories, deepest first
        for d in sorted((p for p in cache_dir.rglob("*") if p.is_dir()), reverse=True):
            if not any(d.iterdir()):
                d.rmdir()

    logger.info("Cleared %d cached files", count)
    return count

cache_info()

Return information about the current cache state.

Returns:

Type Description
dict

dict with keys: enabled, directory, sources, total_files, total_size_mb

Example::

>>> from aeolus.cache import cache_info
>>> info = cache_info()
>>> print(f"Cache: {info['total_files']} files, {info['total_size_mb']:.1f} MB")
Source code in src/aeolus/cache.py
def cache_info() -> dict:
    """
    Return information about the current cache state.

    Returns:
        dict with keys: enabled, directory, sources, total_files, total_size_mb

    Example::

        >>> from aeolus.cache import cache_info
        >>> info = cache_info()
        >>> print(f"Cache: {info['total_files']} files, {info['total_size_mb']:.1f} MB")
    """
    cache_dir = _get_cache_dir()
    # A file may be removed (clear_cache, another process) between listing
    # and stat; skip it rather than crash a read-only introspection call.
    sizes = {}
    for f in cache_dir.rglob("*.parquet"):
        try:
            sizes[f] = f.stat().st_size
        except OSError:
            continue
    total_size = sum(sizes.values())
    entries = _entries_dir()
    current = {f: size for f, size in sizes.items() if entries in f.parents}
    sources = sorted({f.parent.name for f in current})

    return {
        "enabled": _cache_enabled,
        "directory": str(cache_dir),
        "sources": sources,
        "total_files": len(current),
        "total_size_mb": sum(current.values()) / (1024 * 1024),
        # Written by an earlier cache version: never served; clear_cache() removes them
        "legacy_files": len(sizes) - len(current),
    }

is_enabled()

Return whether caching is currently enabled.

Source code in src/aeolus/cache.py
def is_enabled() -> bool:
    """Return whether caching is currently enabled."""
    return _cache_enabled