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
53 changes: 43 additions & 10 deletions pointCollection/io_utils.py
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,16 @@
'NSIDC': 'https://data.nsidc.earthdatacloud.nasa.gov/s3credentials',
}

# The broker call is tried this many times, this far apart. The MAAP API
# stops answering when several hundred jobs start together: on 2026-10-01, 61
# of 556 jobs submitted at once lost their one attempt, each after ~135 s.
MAAP_BROKER_ATTEMPTS = 5
MAAP_BROKER_PAUSE_S = 10

# daac -> why the last broker call failed, for the error get_s3fs() raises
# when the earthaccess fallback has nothing either.
_BROKER_FAILURES = {}

# Block size for remote reads that pull windows out of a large file. fsspec's
# 5 MiB default is sized for reading a file end to end; a windowed read of a
# chunked HDF5 file touches scattered chunks, and the read-ahead is then mostly
Expand Down Expand Up @@ -140,7 +150,19 @@ def get_s3fs(daac='NSIDC', **kwargs):
if fs is None:
import earthaccess
# earthaccess manages its own session; we do not know its expiry.
fs, expires_at = earthaccess.get_s3fs_session(daac=daac, **kwargs), None
try:
fs, expires_at = earthaccess.get_s3fs_session(daac=daac, **kwargs), None
except Exception as exc:
broker = _BROKER_FAILURES.get(daac)
if broker is None:
raise
# On a MAAP DPS worker earthaccess has no login, so the
# fallback fails with an AttributeError about a NoneType that
# says nothing of the cause. Name both.
raise RuntimeError(
f'no {daac} S3 credentials: {broker}; and the earthaccess '
f'fallback has no login here ({type(exc).__name__}: {exc})'
) from exc
_S3FS_CACHE[key] = (fs, expires_at)
return _S3FS_CACHE[key][0]

Expand All @@ -154,6 +176,10 @@ def _s3fs_from_maap(daac, **kwargs):
environment or the broker will not answer, so the caller falls back to
earthaccess. Off MAAP this costs one dict lookup and returns (None, None).

The broker call is tried MAAP_BROKER_ATTEMPTS times, MAAP_BROKER_PAUSE_S
apart, before giving up; the reason it gave up is kept in _BROKER_FAILURES
so that get_s3fs() can report it if earthaccess cannot help either.

This exists because a MAAP DPS worker has NO Earthdata credentials: it runs
as root with no ~/.netrc, and earthaccess's netrc and environment
strategies both come up empty there. What it does have is MAAP's own auth
Expand All @@ -173,8 +199,10 @@ def _s3fs_from_maap(daac, **kwargs):
cause.
"""
import os
import time
import warnings

_BROKER_FAILURES.pop(daac, None)
endpoint = MAAP_S3_CREDENTIALS_ENDPOINTS.get(str(daac).upper())
if endpoint is None:
return None, None
Expand All @@ -190,15 +218,20 @@ def _s3fs_from_maap(daac, **kwargs):
f'falling back to earthaccess for {daac}.')
return None, None

try:
creds = MAAP(
maap_host=os.environ.get('MAAP_API_HOST', 'api.maap-project.org')
).aws.earthdata_s3_credentials(endpoint)
return _s3fs_with_credentials(creds, **kwargs), _expiry_timestamp(creds)
except Exception as exc:
warnings.warn(f'MAAP could not broker {daac} credentials from {endpoint} '
f'({type(exc).__name__}: {exc}); falling back to earthaccess.')
return None, None
for attempt in range(1, MAAP_BROKER_ATTEMPTS + 1):
try:
creds = MAAP(
maap_host=os.environ.get('MAAP_API_HOST', 'api.maap-project.org')
).aws.earthdata_s3_credentials(endpoint)
return _s3fs_with_credentials(creds, **kwargs), _expiry_timestamp(creds)
except Exception as exc:
last = f'{type(exc).__name__}: {exc}'
if attempt < MAAP_BROKER_ATTEMPTS:
time.sleep(MAAP_BROKER_PAUSE_S)
_BROKER_FAILURES[daac] = (f'MAAP could not broker {daac} credentials from {endpoint} '
f'in {MAAP_BROKER_ATTEMPTS} attempts (last error: {last})')
warnings.warn(f'{_BROKER_FAILURES[daac]}; falling back to earthaccess.')
return None, None


# Remembers the outcome of try_earthaccess_login(): None = not yet attempted.
Expand Down
19 changes: 14 additions & 5 deletions pointCollection/ps_scale_for_lat.py
Original file line number Diff line number Diff line change
@@ -1,6 +1,5 @@
import numpy as np
import warnings
warnings.filterwarnings("ignore")

def ps_scale_for_lat(lat):
'''
Expand Down Expand Up @@ -31,7 +30,15 @@ def ps_scale_for_lat(lat):
https://pubs/usgs.gov/pp/1395/report.pdf
'''

if np.nanmean(lat) > 0:
# Warnings this function knows about are silenced HERE, around the lines
# that raise them. A module-level warnings.filterwarnings("ignore") did
# that until 2026-10, and silenced every warning in the importing process
# with it -- including the one saying why a credential request failed.
with warnings.catch_warnings():
# an all-NaN input: "Mean of empty slice"; the result is NaN anyway
warnings.simplefilter("ignore", category=RuntimeWarning)
mean_lat = np.nanmean(lat)
if mean_lat > 0:
hemisphere=1
else:
hemisphere=-1
Expand Down Expand Up @@ -64,9 +71,11 @@ def ps_scale_for_lat(lat):
#print(t)

# distance scaling including special case of the pole
k = t/m *mc_tc
kp = 0.5*mc_tc*np.sqrt(((1.0+e)**(1.0+e))*((1.0-e)**(1.0-e)))
scale = np.where(np.isclose(latr,np.pi/2.0),1.0/kp,1.0/k)
# m is zero at the pole, where kp is used instead of k
with np.errstate(divide='ignore', invalid='ignore'):
k = t/m *mc_tc
kp = 0.5*mc_tc*np.sqrt(((1.0+e)**(1.0+e))*((1.0-e)**(1.0-e)))
scale = np.where(np.isclose(latr,np.pi/2.0),1.0/kp,1.0/k)
return scale
#
# check: at S pole (this comes out right!)
Expand Down
142 changes: 142 additions & 0 deletions tests/test_maap_broker.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,142 @@
"""
The MAAP-brokered DAAC credentials in io_utils.get_s3fs(): the broker call is
retried, and when neither the broker nor earthaccess has credentials the error
says why. No network: stand-ins replace maap-py and earthaccess.

Background (2026-10-01, 556 jobs submitted to MAAP DPS at once): 61 lost their
single broker attempt, fell back to earthaccess -- which has no login on a
worker -- and died with "'NoneType' object has no attribute
'get_s3_filesystem'", the real reason hidden by a process-wide
warnings.filterwarnings("ignore") in ps_scale_for_lat.py.
"""
import subprocess
import sys
import types
import warnings

import numpy as np
import pytest

import pointCollection as pc
from pointCollection import io_utils

CREDS = {'accessKeyId': 'id', 'secretAccessKey': 'secret', 'sessionToken': 'token',
'expiration': '2099-01-01 00:00:00+00:00'}


@pytest.fixture
def broker(monkeypatch):
"""A stand-in maap-py whose broker fails `failures` times, then answers."""
state = {'calls': 0, 'failures': 0, 'naps': []}

class AWS:
def earthdata_s3_credentials(self, endpoint):
state['calls'] += 1
if state['calls'] <= state['failures']:
raise ConnectionError('timed out')
return dict(CREDS)

class MAAP:
def __init__(self, maap_host=None):
self.aws = AWS()

package, module = types.ModuleType('maap'), types.ModuleType('maap.maap')
module.MAAP = MAAP
monkeypatch.setitem(sys.modules, 'maap', package)
monkeypatch.setitem(sys.modules, 'maap.maap', module)
monkeypatch.setenv('MAAP_PGT', 'set')
monkeypatch.setattr('time.sleep', state['naps'].append)
monkeypatch.setattr(io_utils, '_s3fs_with_credentials', lambda creds, **kw: ('fs', creds))
monkeypatch.setattr(io_utils, '_S3FS_CACHE', {})
monkeypatch.setattr(io_utils, '_BROKER_FAILURES', {})
return state


def fake_earthaccess(monkeypatch, session):
module = types.ModuleType('earthaccess')

def get_s3fs_session(daac=None, **kwargs):
if isinstance(session, Exception):
raise session
return session
module.get_s3fs_session = get_s3fs_session
monkeypatch.setitem(sys.modules, 'earthaccess', module)


def test_the_broker_call_is_retried(broker):
broker['failures'] = 2
fs, expires = io_utils._s3fs_from_maap('NSIDC')
assert fs == ('fs', CREDS) and expires is not None
assert broker['calls'] == 3
assert broker['naps'] == [io_utils.MAAP_BROKER_PAUSE_S] * 2


def test_giving_up_warns_with_the_reason(broker):
broker['failures'] = 99
with pytest.warns(UserWarning, match=r'in 5 attempts \(last error: ConnectionError: timed out\)'):
assert io_utils._s3fs_from_maap('NSIDC') == (None, None)
assert broker['calls'] == io_utils.MAAP_BROKER_ATTEMPTS
assert len(broker['naps']) == io_utils.MAAP_BROKER_ATTEMPTS - 1 # none after the last


def test_no_broker_and_no_earthaccess_login_names_both(broker, monkeypatch):
# a DPS worker: this is the AttributeError earthaccess raises with no login
broker['failures'] = 99
fake_earthaccess(monkeypatch, AttributeError("'NoneType' object has no attribute 'get_s3_filesystem'"))
with pytest.warns(UserWarning), pytest.raises(RuntimeError) as err:
io_utils.get_s3fs(daac='NSIDC')
message = str(err.value)
assert 'no NSIDC S3 credentials' in message
assert 'ConnectionError: timed out' in message
assert 'earthaccess fallback has no login here (AttributeError' in message


def test_earthaccess_still_serves_when_the_broker_is_down(broker, monkeypatch):
# the ADE with an Earthdata login: the fallback is still the answer
broker['failures'] = 99
fake_earthaccess(monkeypatch, 'earthaccess session')
with pytest.warns(UserWarning, match='falling back to earthaccess'):
assert io_utils.get_s3fs(daac='NSIDC') == 'earthaccess session'


def test_off_maap_an_earthaccess_error_is_left_alone(monkeypatch):
monkeypatch.delenv('MAAP_PGT', raising=False)
monkeypatch.setattr(io_utils, '_S3FS_CACHE', {})
monkeypatch.setattr(io_utils, '_BROKER_FAILURES', {})
fake_earthaccess(monkeypatch, ValueError('not logged in'))
with pytest.raises(ValueError, match='not logged in'):
io_utils.get_s3fs(daac='NSIDC')


def test_a_later_success_forgets_the_earlier_failure(broker):
broker['failures'] = 99
with pytest.warns(UserWarning):
io_utils._s3fs_from_maap('NSIDC')
assert 'NSIDC' in io_utils._BROKER_FAILURES
broker['calls'], broker['failures'] = 0, 0
io_utils._s3fs_from_maap('NSIDC')
assert 'NSIDC' not in io_utils._BROKER_FAILURES


def test_importing_pointCollection_does_not_silence_warnings():
code = ('import warnings; before = list(warnings.filters); import pointCollection;'
'added = [f for f in warnings.filters if f not in before];'
# a BLANKET ignore: no message pattern, every category
'assert not [f for f in added if f[0] == "ignore" and f[1] is None and f[2] is Warning], added')
result = subprocess.run([sys.executable, '-c', code], capture_output=True, text=True)
assert result.returncode == 0, result.stderr


@pytest.mark.parametrize('lat', [np.array([90., 89.9, 70.]), np.array([-90., -71.]), 90.0,
np.array([np.nan, np.nan])])
def test_ps_scale_for_lat_is_quiet_on_its_own(lat):
with warnings.catch_warnings():
warnings.simplefilter('error')
scale = pc.ps_scale_for_lat(lat)
assert np.shape(scale) == np.shape(lat)


def test_ps_scale_for_lat_values_are_unchanged():
# computed with the module-level filter still in place (main, 42f66cc)
np.testing.assert_allclose(pc.ps_scale_for_lat(np.array([90., 70., 60.])),
[1.03107857, 1.0, 0.96206753], rtol=1e-6)
Loading