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
54 changes: 39 additions & 15 deletions pointCollection/io_utils.py
Original file line number Diff line number Diff line change
Expand Up @@ -30,11 +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
# The broker call is retried after pauses that grow, each jittered by
# +-MAAP_BROKER_JITTER so that jobs which failed together do not retry
# together, and no pause is taken that would end past MAAP_BROKER_BUDGET_S
# from the first try -- a node waiting costs money. 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. The same
# schedule as ATL1415's workspace_credentials.py (plan_pack_tiles.sh K2-K3).
MAAP_BROKER_PAUSES_S = (10, 20, 40, 60, 60) # between tries: up to 6 tries
MAAP_BROKER_JITTER = 0.5
MAAP_BROKER_BUDGET_S = 240

# daac -> why the last broker call failed, for the error get_s3fs() raises
# when the earthaccess fallback has nothing either.
Expand Down Expand Up @@ -176,9 +181,13 @@ 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.
The broker call is tried up to len(MAAP_BROKER_PAUSES_S) + 1 times, after
the pauses above, on ONE MAAP client (building it is itself an API call),
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. HTTP 401
gives up at once: MAAP_PGT was rejected, and retrying the same token
cannot help. There is no per-try timeout: this may run off the main
thread, where SIGALRM is not available.

This exists because a MAAP DPS worker has NO Earthdata credentials: it runs
as root with no ~/.netrc, and earthaccess's netrc and environment
Expand All @@ -199,6 +208,7 @@ def _s3fs_from_maap(daac, **kwargs):
cause.
"""
import os
import random
import time
import warnings

Expand All @@ -218,18 +228,32 @@ def _s3fs_from_maap(daac, **kwargs):
f'falling back to earthaccess for {daac}.')
return None, None

for attempt in range(1, MAAP_BROKER_ATTEMPTS + 1):
client = None
attempts = len(MAAP_BROKER_PAUSES_S) + 1
t0 = time.monotonic()
why = ''
for attempt in range(1, attempts + 1):
try:
creds = MAAP(
maap_host=os.environ.get('MAAP_API_HOST', 'api.maap-project.org')
).aws.earthdata_s3_credentials(endpoint)
if client is None:
client = MAAP(maap_host=os.environ.get('MAAP_API_HOST', 'api.maap-project.org'))
creds = client.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)
if getattr(getattr(exc, 'response', None), 'status_code', None) == 401:
why = ('; HTTP 401: MAAP_PGT was rejected (likely the runner\'s token fetch '
'failed at job start), and a retry cannot help')
break
if attempt == attempts:
break
nap = MAAP_BROKER_PAUSES_S[attempt - 1] * random.uniform(1 - MAAP_BROKER_JITTER,
1 + MAAP_BROKER_JITTER)
if time.monotonic() - t0 + nap > MAAP_BROKER_BUDGET_S:
why = f'; a {nap:.0f} s pause would pass the {MAAP_BROKER_BUDGET_S} s budget'
break
time.sleep(nap)
_BROKER_FAILURES[daac] = (f'MAAP could not broker {daac} credentials from {endpoint} '
f'in {MAAP_BROKER_ATTEMPTS} attempts (last error: {last})')
f'in {attempt} attempts (last error: {last}){why}')
warnings.warn(f'{_BROKER_FAILURES[daac]}; falling back to earthaccess.')
return None, None

Expand Down
60 changes: 52 additions & 8 deletions tests/test_maap_broker.py
Original file line number Diff line number Diff line change
Expand Up @@ -26,26 +26,38 @@

@pytest.fixture
def broker(monkeypatch):
"""A stand-in maap-py whose broker fails `failures` times, then answers."""
state = {'calls': 0, 'failures': 0, 'naps': []}
"""A stand-in maap-py whose broker fails `failures` times (with `error`),
then answers; whose client fails to build `build_failures` times; and a
clock that only sleep() moves, with the jitter fixed at `jitter`."""
state = {'calls': 0, 'failures': 0, 'naps': [], 'builds': 0, 'build_failures': 0,
'error': ConnectionError('timed out'), 'now': 0.0, 'jitter': 1.0}

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

class MAAP:
def __init__(self, maap_host=None):
state['builds'] += 1
if state['builds'] <= state['build_failures']:
raise ConnectionError('config timed out')
self.aws = AWS()

def sleep(s):
state['naps'].append(s)
state['now'] += s

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('time.sleep', sleep)
monkeypatch.setattr('time.monotonic', lambda: state['now'])
monkeypatch.setattr('random.uniform', lambda a, b: state['jitter'] * (a + b) / 2)
monkeypatch.setattr(io_utils, '_s3fs_with_credentials', lambda creds, **kw: ('fs', creds))
monkeypatch.setattr(io_utils, '_S3FS_CACHE', {})
monkeypatch.setattr(io_utils, '_BROKER_FAILURES', {})
Expand All @@ -68,15 +80,47 @@ def test_the_broker_call_is_retried(broker):
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
assert broker['naps'] == list(io_utils.MAAP_BROKER_PAUSES_S[:2])
assert broker['builds'] == 1 # one client for every try


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\)'):
with pytest.warns(UserWarning, match=r'in 6 attempts \(last error: ConnectionError: timed out\)'):
assert io_utils._s3fs_from_maap('NSIDC') == (None, None)
assert broker['calls'] == len(io_utils.MAAP_BROKER_PAUSES_S) + 1
assert broker['naps'] == list(io_utils.MAAP_BROKER_PAUSES_S) # none after the last
assert sum(broker['naps']) <= io_utils.MAAP_BROKER_BUDGET_S


def test_pauses_are_jittered(broker):
broker['failures'], broker['jitter'] = 1, 0.5
io_utils._s3fs_from_maap('NSIDC')
assert broker['naps'] == [io_utils.MAAP_BROKER_PAUSES_S[0] * 0.5]


def test_no_pause_past_the_budget(broker):
# jitter at its top: 15 + 30 + 60 + 90 = 195 s, and the next 90 s would end at 285
broker['failures'], broker['jitter'] = 99, 1.5
with pytest.warns(UserWarning, match=r'in 5 attempts .*a 90 s pause would pass the 240 s budget'):
assert io_utils._s3fs_from_maap('NSIDC') == (None, None)
assert broker['naps'] == [15, 30, 60, 90]


def test_401_stops_at_once(broker):
class HTTPError(Exception):
response = types.SimpleNamespace(status_code=401)
broker['failures'], broker['error'] = 99, HTTPError('401 Unauthorized')
with pytest.warns(UserWarning, match=r'in 1 attempts .*HTTP 401: MAAP_PGT was rejected'):
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
assert broker['calls'] == 1 and broker['naps'] == []


def test_a_client_that_failed_to_build_is_built_again(broker):
broker['build_failures'] = 2
fs, _ = io_utils._s3fs_from_maap('NSIDC')
assert fs == ('fs', CREDS)
assert broker['builds'] == 3 and broker['calls'] == 1


def test_no_broker_and_no_earthaccess_login_names_both(broker, monkeypatch):
Expand Down
Loading