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: 5 additions & 2 deletions .github/workflows/pytests-upstream.yml
Original file line number Diff line number Diff line change
Expand Up @@ -93,6 +93,7 @@ jobs:
echo "CONDA_ENV_FILE=ci/requirements/py${{matrix.python-version}}-core-free.yml" >> $GITHUB_ENV
echo "PYTHON_VERSION=${{ matrix.python-version }}" >> $GITHUB_ENV
echo "LOG_FILE=argopy-tests-Core-Free-Py${{matrix.python-version}}-${{matrix.os}}.log" >> $GITHUB_ENV
echo "PYTHONTRACEMALLOC=20" >> $GITHUB_ENV

- name: Increase swapfile
if: ${{matrix.os == 'ubuntu-latest'}}
Expand All @@ -118,7 +119,7 @@ jobs:

- name: Install argopy
run: |
python -m pip install --no-deps -e .
python -m pip install --no-deps -e .

- name: Version info
run: |
Expand Down Expand Up @@ -261,6 +262,7 @@ jobs:
run: |
echo "CONDA_ENV_FILE=ci/requirements/py${{matrix.python-version}}-all-free.yml" >> $GITHUB_ENV
echo "PYTHON_VERSION=${{ matrix.python-version }}" >> $GITHUB_ENV
echo "PYTHONTRACEMALLOC=20" >> $GITHUB_ENV

- name: Increase swapfile
if: ${{matrix.os == 'ubuntu-latest'}}
Expand All @@ -287,7 +289,7 @@ jobs:

- name: Install argopy
run: |
python -m pip install --no-deps -e .
python -m pip install --no-deps -e .

- name: Version info
run: |
Expand All @@ -313,6 +315,7 @@ jobs:

- name: Test with pytest
run: |
PYTHONTRACEMALLOC=1
python ci/resource_summary.py pytest -ra -v -s -c argopy/tests/pytest.ini --durations=10 \
--report-log output-${{ matrix.python-version }}-log.jsonl

Expand Down
8 changes: 4 additions & 4 deletions argopy/data_fetchers/gdac_data.py
Original file line number Diff line number Diff line change
Expand Up @@ -276,10 +276,10 @@ def to_xarray(
and not self.parallelize
and self.parallel_method == "sequential"
):
warnings.warn(
"Found more than 50 files to load, this may take a while to process sequentially ! "
"Consider using another data source (eg: 'erddap') or the 'parallel=True' option to improve processing time."
)
msg = f"Found more than 50 files to load ({len(URI)} !), this may take a while to process sequentially ! Consider using another data source (eg: 'erddap') or the 'parallel=True' option to improve processing time."
warnings.warn(msg)
log.info(msg)

elif len(URI) == 0:
raise DataNotFound("No data found for: %s" % self.indexfs.cname)

Expand Down
6 changes: 3 additions & 3 deletions argopy/extensions/optical_modeling.py
Original file line number Diff line number Diff line change
Expand Up @@ -104,7 +104,7 @@ def Zeu(
But the euphotic depth can also be estimated using the exponential decay of light with depth, described by Beer's Law [1]_:

.. math::
I(z) = I_0 \exp(-K_{PAR}\\,z)
I(z) = I_0 \\exp(-K_{PAR}\\,z)

If we solve for $I(Z_e)=0.01 I_0$ we get:

Expand Down Expand Up @@ -288,11 +288,11 @@ def Z_iPAR_threshold(

Notes
-----
This is the closest level $z$ in the vertical axis for which PAR is about a threshold value $t$, with some tolerance $\epsilon$:
This is the closest level $z$ in the vertical axis for which PAR is about a threshold value $t$, with some tolerance $\\epsilon$:

.. math::

z | abs(PAR(z) - t) < \epsilon
z | abs(PAR(z) - t) < \\epsilon

A default value of 15 is used because it is the theoretical value below which the Fchla is no longer
quenched (For correction of NPQ purposes).
Expand Down
6 changes: 3 additions & 3 deletions argopy/stores/float/spec.py
Original file line number Diff line number Diff line change
Expand Up @@ -426,7 +426,7 @@ def open_dataset(
cast: bool, optional, default = True
Determine if the dataset variables should be cast or not. This is similar to opening the dataset directly with :class:`xarray.open_dataset` using the ``engine=`argo``` option.
This will be ignored if the ``netCDF4` kwarg is set to True.
\**kwargs
**kwargs
All the other arguments are passed to the GDAC store `open_dataset` method.

Returns
Expand Down Expand Up @@ -461,7 +461,7 @@ def dataset(self, name: str = "prof", **kwargs) -> xr.Dataset:
----------
name: str, optional, default = "prof"
Name of the dataset to open. It can be any key from the dictionary returned by :class:`ArgoFloat.ls_datasets`.
\**kwargs
**kwargs
All the other arguments are passed to the :meth:`ArgoFloat.open_dataset` method.

Returns
Expand Down Expand Up @@ -804,7 +804,7 @@ def profile(self, name: str, **kwargs) -> xr.Dataset:
----------
name: str
Name of the profile file to open. It can be any key from the dictionary returned by :class:`ArgoFloat.ls_profiles`.
\**kwargs
**kwargs
All the other arguments are passed to the :meth:`ArgoFloat.open_profile` method.

Returns
Expand Down
10 changes: 9 additions & 1 deletion argopy/stores/implementations/ftp.py
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@
import warnings
from typing import Literal
from netCDF4 import Dataset
import threading

from ...errors import InvalidMethod, DataNotFound
from ...utils.transformers import drop_variables_not_in_all_datasets
Expand All @@ -18,6 +19,9 @@

log = logging.getLogger("argopy.stores.implementation.ftp")

_cache_lock = threading.Lock()
# Used to lock threads to prevent race to the cached meda-data of fsspec


class ftpstore(httpstore):
"""Argo ftp file system
Expand Down Expand Up @@ -110,7 +114,8 @@ def load_in_memory(url, errors="raise", xr_opts={}):

try:
this_url = self.fs._strip_protocol(url)
data = self.fs.cat_file(this_url)
with _cache_lock:
data = self.fs.cat_file(this_url)
if data is None:
if errors == "raise":
raise DataNotFound(url)
Expand Down Expand Up @@ -207,6 +212,9 @@ def load_lazily(url, errors="raise", xr_opts={}, akoverwrite: bool = False):
if target is not None:
if not netCDF4:
ds = xr.open_dataset(target, **xr_opts)
if not lazy:
ds = ds.load() # materialize into plain numpy arrays, detach from the backend buffer
ds.close() # explicitly release the netCDF4/HDF5 handle right away

if "source" not in ds.encoding:
if isinstance(url, str):
Expand Down
17 changes: 12 additions & 5 deletions argopy/stores/implementations/http.py
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@
from functools import lru_cache
from netCDF4 import Dataset
from urllib.parse import urlparse
import threading

from ...errors import InvalidMethod, DataNotFound
from ...utils import Registry, UriCName
Expand All @@ -33,6 +34,9 @@

log = logging.getLogger("argopy.stores.implementation.http")

_cache_lock = threading.Lock()
# Used to lock threads to prevent race to the cached meda-data of fsspec


class httpstore(ArgoStoreProto):
"""Argo http file system
Expand Down Expand Up @@ -118,7 +122,8 @@ def make_request(
data = None
if n_attempt <= max_attempt:
try:
data = ffs.cat_file(url, **cat_opts)
with _cache_lock:
data = ffs.cat_file(url, **cat_opts)
except FileNotFoundError as e:
if errors == "raise":
raise e
Expand Down Expand Up @@ -369,7 +374,8 @@ def load_lazily(
if not netCDF4:
ds = xr.open_dataset(target, **xr_opts)
if not lazy:
ds = ds.load()
ds = ds.load() # materialize into plain numpy arrays, detach from the backend buffer
ds.close() # explicitly release the netCDF4/HDF5 handle right away

if "source" not in ds.encoding:
if isinstance(url, str):
Expand Down Expand Up @@ -905,9 +911,10 @@ def read_csv(self, url, **kwargs):

"""
url = self.curateurl(url)
# log.debug("Opening/reading csv from: %s" % url)
with self.open(url) as of:
df = pd.read_csv(of, **kwargs)

with _cache_lock:
with self.open(url) as of:
df = pd.read_csv(of, **kwargs)

self.register(url)
return df
Expand Down
21 changes: 16 additions & 5 deletions argopy/stores/implementations/local.py
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@
from pathlib import Path
import warnings
from netCDF4 import Dataset
import threading

from ...options import OPTIONS
from ...errors import InvalidMethod, DataNotFound
Expand All @@ -19,8 +20,12 @@
from ..filesystems import has_distributed, distributed
from ..filesystems import tqdm


log = logging.getLogger("argopy.stores.implementation.local")

_cache_lock = threading.Lock()
# Used to lock threads to prevent race to the cached meda-data of fsspec


class filestore(ArgoStoreProto):
"""Argo local file system
Expand Down Expand Up @@ -71,8 +76,9 @@ def open_json(self, url, errors: Literal['raise', 'silent', 'ignore'] = 'raise',
if "js_opts" in kwargs:
js_opts.update(kwargs["js_opts"])

with self.open(url, **open_opts) as of:
js = json.load(of, **js_opts)
with _cache_lock:
with self.open(url, **open_opts) as of:
js = json.load(of, **js_opts)

if len(js) == 0:
if errors == "raise":
Expand Down Expand Up @@ -121,7 +127,8 @@ def load_in_memory(path, errors="raise", xr_opts={}):
tuple: (data, _) or (None, _) if errors == "ignore"
"""
try:
data = self.fs.cat_file(path)
with _cache_lock:
data = self.fs.cat_file(path)

if data[0:3] != b"CDF" and data[0:3] != b"\x89HD":
raise TypeError(
Expand Down Expand Up @@ -198,6 +205,9 @@ def load_lazily(path, errors="raise", xr_opts={}, akoverwrite: bool = False):
if target is not None:
if not netCDF4:
ds = xr.open_dataset(target, **xr_opts)
if not lazy:
ds = ds.load() # materialize into plain numpy arrays, detach from the backend buffer
ds.close() # explicitly release the netCDF4/HDF5 handle right away

if "source" not in ds.encoding:
if isinstance(path, str):
Expand Down Expand Up @@ -426,6 +436,7 @@ def read_csv(self, path, **kwargs):
:class:`pandas.DataFrame`
"""
log.debug("Reading csv: %s" % path)
with self.open(path) as of:
df = pd.read_csv(of, **kwargs)
with _cache_lock:
with self.open(path) as of:
df = pd.read_csv(of, **kwargs)
return df
38 changes: 35 additions & 3 deletions argopy/tests/conftest.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,10 +2,12 @@
import os
import logging
import shutil
import pytest
import threading

sys.path.append(os.path.join(os.path.dirname(__file__), 'helpers'))
from mocked_ftp import mocked_ftpserver
from mocked_http import mocked_httpserver
from argopy.tests.helpers.mocked_ftp import mocked_ftpserver
from argopy.tests.helpers.mocked_http import mocked_httpserver


log = logging.getLogger("argopy.tests.conftests")
Expand All @@ -24,4 +26,34 @@ def pytest_sessionfinish(session, exitstatus):
pass
log.debug("Ending tests session")
log.debug("Final session state: %s" % session)
pass
pass

@pytest.fixture(autouse=True)
def _resource_tracker(request):
"""Track resource usage around every test to catch what leaks."""

proc = f"/proc/{os.getpid()}"

def _fds():
try:
return len(os.listdir(f"{proc}/fd"))
except Exception:
return -1

def _threads():
return threading.active_count()

before = dict(fds=_fds(), threads=_threads())
yield
after = dict(fds=_fds(), threads=_threads())

delta_fds = after["fds"] - before["fds"]
delta_threads = after["threads"] - before["threads"]

# Only log when something looks wrong
if delta_fds > 10 or delta_threads > 2:
print(
f"\n[RESOURCE LEAK] {request.node.nodeid}\n"
f" FDs: {before['fds']} : {after['fds']} (diff:{delta_fds:+d})\n"
f" Threads: {before['threads']} : {after['threads']} (diff:{delta_threads:+d})\n"
)
58 changes: 21 additions & 37 deletions argopy/tests/helpers/mocked_http.py
Original file line number Diff line number Diff line change
Expand Up @@ -29,7 +29,7 @@
from pathlib import Path
import threading
from collections import ChainMap
from http.server import BaseHTTPRequestHandler, HTTPServer
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
import pytest
import logging
from urllib.parse import unquote
Expand All @@ -44,25 +44,6 @@
False # Should we list all files/uris available from the mocked server in the log ?
)

import socket


def _free_port() -> int:
"""Return a free port number on localhost.

bind("127.0.0.1", 0) asks the OS to assign a free port, which we read back
with getsockname"""
s = socket.socket()
s.bind(("127.0.0.1", 0))
p = s.getsockname()[1]
s.close()
return p


port = _free_port()
mocked_server_address = "http://127.0.0.1:%i" % port


"""
Load test data and create a dictionary mapping of URL requests as keys, and expected responses as values

Expand Down Expand Up @@ -174,22 +155,23 @@ def __init__(self, *args, **kwargs):
def _respond(self, code=200, headers=None, data=b""):
headers = headers or {}
headers.update({"User-Agent": "Mocked http server for unit tests"})
self.send_response(code)
for k, v in headers.items():
self.send_header(k, str(v))
self.end_headers()
if data:
if not isinstance(data, (bytes, bytearray)):
data = _read(data)
try:
try:
self.send_response(code)
for k, v in headers.items():
self.send_header(k, str(v))
Comment thread
gmaze marked this conversation as resolved.
Dismissed
self.end_headers()
if data:
if not isinstance(data, (bytes, bytearray)):
data = _read(data)
self.wfile.write(data)
except socket.error as e:
# socket error [Errno 32] Broken pipe
# This might be happening when a client program doesn't wait till all the data from the server is
# received and simply closes a socket
if "32" not in str(e):
log.debug("socket error %s" % str(e))
pass
except (BrokenPipeError, ConnectionResetError, socket.error) as e:
# socket error [Errno 32] Broken pipe
# This might be happening when a client program doesn't wait till all the data from the server is
# received and simply closes a socket
log.debug("socket error while responding: %s" % str(e))
# if "32" not in str(e):
# log.debug("socket error %s" % str(e))
# pass

def log_message(self, format, *args):
# Quiet logging !
Expand Down Expand Up @@ -301,8 +283,10 @@ def do_HEAD(self):

@contextlib.contextmanager
def serve_mocked_httpserver():
server_address = ("", port)
httpd = HTTPServer(server_address, HTTPTestHandler)
httpd = ThreadingHTTPServer(("127.0.0.1", 0), HTTPTestHandler)
port = httpd.server_address[1]
mocked_server_address = "http://127.0.0.1:%i" % port

th = threading.Thread(target=httpd.serve_forever)
th.daemon = True
th.start()
Expand Down
Loading
Loading