From 7bc60dcaeaa7383f5abd027c0d511415101abbf4 Mon Sep 17 00:00:00 2001 From: Mike McCann Date: Thu, 26 Mar 2026 11:01:29 -0700 Subject: [PATCH 1/6] fix(create_products): skip 2column plot when no data variables exist Add early return in plot_2column() to avoid creating an empty _2column_cmocean.png when no instrument variables are present. --- .vscode/launch.json | 4 +++- src/data/create_products.py | 9 +++++++++ 2 files changed, 12 insertions(+), 1 deletion(-) diff --git a/.vscode/launch.json b/.vscode/launch.json index c3c725b..ddd6b88 100644 --- a/.vscode/launch.json +++ b/.vscode/launch.json @@ -460,7 +460,9 @@ // Test sipper_odv() with a log_file that has an ESP Sample //"args": ["-v", "1", "--log_file", "daphne/missionlogs/2026/20260316_20260318/20260317T191958/202603171919_202603181628.nc4", "--no_cleanup"] // Test sipper_odv() with a log_file that has an ESP Sample - "args": ["-v", "1", "--log_file", "ahi/missionlogs/2025/20250414_20250418/20250415T040019/202504150400_202504152346.nc4", "--no_cleanup"] + //"args": ["-v", "1", "--log_file", "ahi/missionlogs/2025/20250414_20250418/20250415T040019/202504150400_202504152346.nc4", "--no_cleanup"] + // Test not writing _2column_cmocean.png file if there is no data to plot + "args": ["-v", "1", "--log_file", "ahi/missionlogs/2025/20250128_20250131/20250131T051404/202501310514_202501310535.nc4", "--clobber"] }, ] diff --git a/src/data/create_products.py b/src/data/create_products.py index 6887817..d1c5e4f 100755 --- a/src/data/create_products.py +++ b/src/data/create_products.py @@ -1638,6 +1638,15 @@ def plot_2column(self) -> str: # noqa: C901, PLR0912, PLR0915 self._open_ds() + # Early return if no plot variables present in dataset + # Use a quick pre-check with LRAUV or Dorado variables (excluding computed 'density') + plot_variables = self._get_plot_variables(None if self._is_lrauv() else "ctd1") + if not any(var in self.ds for var, _ in plot_variables if var != "density"): + self.logger.warning( + "No plot variables found in dataset, skipping plot_2column", + ) + return None + idist, iz, distnav = self._grid_dims() if idist.size == 0 or iz.size == 0 or distnav.size == 0: self.logger.warning("Skipping plot_2column due to missing gridding dimensions") From a00e5013c8281d97cef07adc25557628ee67bd91 Mon Sep 17 00:00:00 2001 From: Mike McCann Date: Wed, 1 Apr 2026 14:46:41 -0700 Subject: [PATCH 2/6] Add instructions for creating a .env file The .env file is meant to hold an API token for updating the SSDS_Metadata database with ProcessRun information. --- .env.example | 18 ++++++++++++++++ .vscode/launch.json | 7 ++++-- README.md | 52 +++++++++++++++++++++++++++------------------ docker-compose.yml | 1 + 4 files changed, 55 insertions(+), 23 deletions(-) create mode 100644 .env.example diff --git a/.env.example b/.env.example new file mode 100644 index 0000000..e5aa97a --- /dev/null +++ b/.env.example @@ -0,0 +1,18 @@ +# Copy this file to .env and fill in values for your host/container. +# .env is ignored by git. + +# Existing docker-compose volume vars +M3_VOL=/Volumes/M3 +AUVCTD_VOL=/Volumes/AUVCTD +LRAUV_VOL=/Volumes/LRAUV +CALIBRATION_VOL=/Volumes/DMO +WORK_VOL=/opt/docker_auv-python_vols/data +HOST_NAME=example.shore.mbari.org +GMT_LIBRARY_PATH=/usr/lib/x86_64-linux-gnu/ + +# SSDS API endpoint (set to your dev server as needed) +SSDS_API_BASE=https://mooring-ssds.shore.mbari.org/api + +# SSDS provenance auth stubs (API key) +SSDS_API_KEY= +SSDS_API_KEY_HEADER=X-API-Key diff --git a/.vscode/launch.json b/.vscode/launch.json index ddd6b88..98a94ff 100644 --- a/.vscode/launch.json +++ b/.vscode/launch.json @@ -394,6 +394,7 @@ "request": "launch", "program": "${workspaceFolder}/src/data/process_lrauv.py", "console": "integratedTerminal", + "envFile": "${workspaceFolder}/.env", // Lots bad time values in brizo 20250914T080941 due to memory corruption on the vehicle //"args": ["-v", "1", "--log_file", "brizo/missionlogs/2025/20250909_20250915/20250914T080941/202509140809_202509150109.nc4"] //"args": ["-v", "2", "--log_file", "brizo/missionlogs/2025/20250909_20250915/20250914T080941/202509140809_202509150109.nc4", "--clobber"] @@ -459,10 +460,12 @@ //"args": ["-v", "1", "--log_file", "pontus/missionlogs/2025/20250107_20250123/20250111T233141/202501112331_202501131744.nc4", "--no_cleanup"] // Test sipper_odv() with a log_file that has an ESP Sample //"args": ["-v", "1", "--log_file", "daphne/missionlogs/2026/20260316_20260318/20260317T191958/202603171919_202603181628.nc4", "--no_cleanup"] - // Test sipper_odv() with a log_file that has an ESP Sample + // Test ahi mission that has Backseat Planktivore data //"args": ["-v", "1", "--log_file", "ahi/missionlogs/2025/20250414_20250418/20250415T040019/202504150400_202504152346.nc4", "--no_cleanup"] // Test not writing _2column_cmocean.png file if there is no data to plot - "args": ["-v", "1", "--log_file", "ahi/missionlogs/2025/20250128_20250131/20250131T051404/202501310514_202501310535.nc4", "--clobber"] + //"args": ["-v", "1", "--log_file", "ahi/missionlogs/2025/20250128_20250131/20250131T051404/202501310514_202501310535.nc4", "--clobber"] + // Test --update_ssds_provenance with a short log_file + "args": ["-v", "1", "--log_file", "ahi/missionlogs/2025/20250128_20250131/20250131T051404/202501310514_202501310535.nc4", "--update_ssds_provenance"] }, ] diff --git a/README.md b/README.md index 9106816..b0a435a 100644 --- a/README.md +++ b/README.md @@ -103,28 +103,38 @@ First time use with Docker on a server using a service account: * cd /opt # There should be an `auv-python` directory here that is writable by docker_user * git clone git@github.com:mbari-org/auv-python.git * cd auv-python -* Create a .env file in `/opt/auv-python` with the following contents: - `M3_VOL=` - `AUVCTD_VOL=` - `LRAUV_VOL=` - `CALIBRATION_VOL=` - `WORK_VOL=/data` - `HOST_NAME=` - `GMT_LIBRARY_PATH=/usr/lib/x86_64-linux-gnu/` +* Create a .env file in `/opt/auv-python` using .env.example as a guide + After installation and when logging into the server again mission data can be processed thusly: -* Setting up environment and printing help message: - `sudo -u docker_user -i` - `cd /opt/auv-python` - `git pull` # To get new changes, e.g. mission added to src/data/dorado_info.py - `export DOCKER_USER_ID=$(id -u)` - `docker compose build` - `docker compose run --rm auvpython src/data/process_i2map.py --help` -* To actually process a mission and have the processed data copied to the archive use the `-v` and `--clobber` options, e.g.: - `docker compose run --rm auvpython src/data/process_dorado.py --mission 2025.139.04 -v --clobber --noinput` -* To process LRAUV data for a specific vehicle and time range: - `docker compose run --rm auvpython src/data/process_lrauv.py --auv_name tethys --start 20250401T000000 --end 20250502T000000 -v --noinput` -* To process a specific LRAUV log file: - `docker compose run --rm auvpython src/data/process_lrauv.py --log_file tethys/missionlogs/2012/20120908_20120920/20120917T025522/201209170255_201209171110.nc4 -v --noinput` +* Setting up environment and printing help message: +``` + sudo -u docker_user -i + cd /opt/auv-python + git pull + export DOCKER_USER_ID=$(id -u) + docker compose build + docker compose run --rm auvpython src/data/process_i2map.py --help +``` + +* To actually process a mission and have the processed data copied to the archive use the `-v` and `--clobber` options, e.g.: +``` + docker compose run --rm auvpython src/data/process_dorado.py --mission 2025.139.04 -v --clobber --noinput +``` + +* To process LRAUV data for a specific vehicle and time range: +``` + docker compose run --rm auvpython src/data/process_lrauv.py --auv_name tethys --start 20250401T000000 --end 20250502T000000 -v --clobber --noinput +``` + +* For missions/log_files to be processed sequentially add the `--num_cores 1` option: +``` + docker compose run --rm auvpython src/data/process_lrauv.py --start 20250101T000000 --end 20260101T000000 -v --clobber --noinput --num_cores 1 +``` + +* To process a specific LRAUV log file: +``` + docker compose run --rm auvpython src/data/process_lrauv.py --log_file tethys/missionlogs/2012/20120908_20120920/20120917T025522/201209170255_201209171110.nc4 -v --clobber --noinput +``` -- diff --git a/docker-compose.yml b/docker-compose.yml index 746348b..7b08125 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -1,3 +1,4 @@ +# Create .env from .env.example, then set values for your environment. # Env variables required in .env, e.g. for Mac: # M3_VOL=/Volumes/M3 # AUVCTD_VOL=/Volumes/AUVCTD From 20aa2ce38415b70ce960804cea63c1c1cea85c31 Mon Sep 17 00:00:00 2001 From: Mike McCann Date: Wed, 1 Apr 2026 14:48:16 -0700 Subject: [PATCH 3/6] Add --update_ssds_provenance option This uses a new provenance.py module for submitting ProcessRun data set provenance information using the CRUD-style REST API of the mooring-ssds project. --- pyproject.toml | 1 + src/data/common_args.py | 5 + src/data/process.py | 85 ++++++++ src/data/provenance.py | 463 ++++++++++++++++++++++++++++++++++++++++ 4 files changed, 554 insertions(+) create mode 100644 src/data/provenance.py diff --git a/pyproject.toml b/pyproject.toml index 04dd4ba..aec52fc 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -28,6 +28,7 @@ dependencies = [ "pygmt==0.16", "pyproj>=3.7.1", "pysolar>=0.13", + "requests>=2.31.0", "rolling>=0.5.0", "seawater>=3.3.5", "statsmodels>=0.14.4", diff --git a/src/data/common_args.py b/src/data/common_args.py index dcfd242..12cc853 100644 --- a/src/data/common_args.py +++ b/src/data/common_args.py @@ -87,6 +87,11 @@ def get_processing_parser(): action="store_true", help="Don't re-process existing output files", ) + parser.add_argument( + "--update_ssds_provenance", + action="store_true", + help="Submit/update provenance records in the SSDS_Metadata database", + ) return parser diff --git a/src/data/process.py b/src/data/process.py index 572d61c..e8b6e64 100755 --- a/src/data/process.py +++ b/src/data/process.py @@ -75,6 +75,7 @@ class data are: download_process and calibrate, while for LRAUV class data from logs2netcdfs import BASE_PATH, MISSIONLOGS, MISSIONNETCDFS, AUV_NetCDF from lopcToNetCDF import LOPC_Processor, UnexpectedAreaOfCode from nc42netcdfs import BASE_LRAUV_PATH, BASE_LRAUV_WEB, Extract +from provenance import get_dods_url, submit_process_run from resample import ( AUVCTD_OPENDAP_BASE, FLASH_THRESHOLD, @@ -177,6 +178,7 @@ def __init__(self, auv_name, vehicle_dir, mount_dir, calibration_dir, config=Non "no_cleanup": False, "skip_download_process": False, "archive_only_products": False, + "update_ssds_provenance": False, "num_cores": None, # Filtering/processing params (only used in from_args, not common_config) "start_year": None, @@ -708,6 +710,34 @@ def create_products(self, mission: str = None, log_file: str = None) -> None: cp.sipper_odv() cp.logger.removeHandler(self.log_handler) + def _submit_provenance( # noqa: PLR0913 + self, + output_nc: str, + input_files: list[str], + pr_start: str, + pr_end: str, + script_name: str = "src/data/process.py", + log_file: str | None = None, + ) -> None: + """Submit a provenance record — failures are logged, never raised.""" + try: + if not Path(output_nc).exists(): + self.logger.debug("Output %s not found, skipping provenance", output_nc) + return + log_url = get_dods_url(log_file) if log_file else None + submit_process_run( + nc_file_path=output_nc, + input_uris=input_files, + pr_start=pr_start, + pr_end=pr_end, + script_name=script_name, + cmd_line_args=self.commandline, + log_file_url=log_url, + log=self.logger, + ) + except Exception: # noqa: BLE001 + self.logger.warning("Provenance submission failed", exc_info=True) + def email(self, mission: str) -> None: self.logger.info("Sending notification email for %s", mission) email = Emailer() @@ -771,6 +801,7 @@ def cleanup(self, mission: str = None, log_file: str = None) -> None: self.logger.error("Either mission or log_file must be provided for cleanup.") def process_mission(self, mission: str, src_dir: str = "") -> None: # noqa: C901, PLR0912, PLR0915 + _pr_start = datetime.now(tz=UTC).isoformat() netcdfs_dir = Path( self.config["base_path"], self.auv_name, @@ -854,6 +885,39 @@ def process_mission(self, mission: str, src_dir: str = "") -> None: # noqa: C90 self.align(mission) self.resample(mission) self.create_products(mission) + if self.config["update_ssds_provenance"]: + resampled = Path( + self.config["base_path"], + self.auv_name, + MISSIONNETCDFS, + mission, + f"{self.auv_name}_{mission}_{FREQ}.nc", + ) + self._submit_provenance( + output_nc=str(resampled), + input_files=[ + str( + Path( + self.config["base_path"], + self.auv_name, + MISSIONLOGS, + mission, + ) + ) + ], + pr_start=_pr_start, + pr_end=datetime.now(tz=UTC).isoformat(), + script_name="src/data/process_dorado.py", + log_file=str( + Path( + self.config["base_path"], + self.auv_name, + MISSIONNETCDFS, + mission, + f"{self.auv_name}_{mission}_{LOG_NAME}", + ) + ), + ) # self.archive() is called in finally: blocks in process_missions() def process_mission_job(self, mission: str, src_dir: str = "") -> None: @@ -1041,6 +1105,7 @@ def combine(self, log_file: str) -> None: @log_file_processor def process_log_file(self, log_file: str) -> None: + _pr_start = datetime.now(tz=UTC).isoformat() netcdfs_dir = Path(BASE_LRAUV_PATH, Path(log_file).parent) Path(netcdfs_dir).mkdir(parents=True, exist_ok=True) self.log_handler = logging.FileHandler( @@ -1074,6 +1139,21 @@ def process_log_file(self, log_file: str) -> None: ) if resampled_file.exists(): self.create_products(log_file=log_file) + if self.config["update_ssds_provenance"]: + self._submit_provenance( + output_nc=str(resampled_file), + input_files=[get_dods_url(str(Path(BASE_LRAUV_PATH, log_file)))], + pr_start=_pr_start, + pr_end=datetime.now(tz=UTC).isoformat(), + script_name="src/data/process_lrauv.py", + log_file=str( + Path( + BASE_LRAUV_PATH, + Path(log_file).parent, + f"{Path(log_file).stem}_processing.log", + ) + ), + ) else: self.logger.warning( "Resampled file %s not found, skipping create_products step", @@ -1340,6 +1420,11 @@ def process_command_line(self): "and append to the netCDF file name" ), ) + parser.add_argument( + "--update_ssds_provenance", + action="store_true", + help="Submit/update provenance records in the SSDS_Metadata database", + ) parser.add_argument( "--num_cores", action="store", diff --git a/src/data/provenance.py b/src/data/provenance.py new file mode 100644 index 0000000..1a33bd2 --- /dev/null +++ b/src/data/provenance.py @@ -0,0 +1,463 @@ +"""Record data processing provenance in the SSDS_Metadata database. + +Replaces the legacy Perl ``submitDStoNetCDFProcessRun`` with authenticated +REST API calls to the mooring-ssds service. Every public helper is +designed to be called from ``process.py`` after a successful processing +run, or standalone via the CLI at the bottom of this file. + +See: https://github.com/mbari-org/auv-python/issues/144 +""" + +import argparse +import logging +import os +import sys +from datetime import UTC, datetime +from pathlib import Path +from socket import gethostname + +import git +import requests + +logger = logging.getLogger(__name__) + +# --------------------------------------------------------------------------- +# Constants +# --------------------------------------------------------------------------- +ENV_SSDS_API_BASE = "SSDS_API_BASE" +SSDS_API_BASE = os.environ.get(ENV_SSDS_API_BASE, "https://mooring-ssds.shore.mbari.org/api") +GIT_WEB_BASE = "https://github.com/mbari-org/auv-python/blob" +GITHUB_REPO = "https://github.com/mbari-org/auv-python" +REQUEST_TIMEOUT = 30 + +# Environment variable names used for authentication stubs. +ENV_SSDS_API_KEY = "SSDS_API_KEY" # noqa: S105 +ENV_SSDS_API_KEY_HEADER = "SSDS_API_KEY_HEADER" # noqa: S105 + +# Maps local absolute path prefixes to OPeNDAP URL prefixes. +# Mirrors the Perl ``%web_lookup`` hash in ssds_util.pl. +# Computed from the project layout, matching BASE_LRAUV_PATH / BASE_PATH. +_PROJECT_DATA = Path(__file__).resolve().parent.parent.parent / "data" +PATH_TO_URL_MAP: dict[str, str] = { + str(_PROJECT_DATA / "auv_data"): "http://dods.mbari.org/opendap/data/auvctd", + str(_PROJECT_DATA / "lrauv_data"): "http://dods.mbari.org/opendap/data/lrauv", +} + + +# --------------------------------------------------------------------------- +# Helpers +# --------------------------------------------------------------------------- +def _get_git_version() -> str: + """Return the current git commit hash, or ``'unknown'``.""" + try: + repo = git.Repo(Path(__file__).resolve().parent, search_parent_directories=True) + except Exception: # noqa: BLE001 + logger.debug("Could not determine git version", exc_info=True) + return "unknown" + else: + return repo.head.commit.hexsha + + +def build_authenticated_session( + api_key: str | None = None, + api_key_header: str | None = None, +) -> requests.Session: + """Return a requests.Session with auth headers from args or env. + + Resolution order: + 1) Explicit function args + 2) Environment variables + + Supported auth stubs: + - API key via configurable header (default: X-API-Key) + """ + session = requests.Session() + + resolved_api_key = api_key or os.environ.get(ENV_SSDS_API_KEY) + resolved_api_key_header = ( + api_key_header or os.environ.get(ENV_SSDS_API_KEY_HEADER) or "X-API-Key" + ) + + if resolved_api_key: + session.headers[resolved_api_key_header] = resolved_api_key + + return session + + +def get_dods_url(nc_file_path: str) -> str: + """Translate a local NetCDF path to its OPeNDAP URL. + + Walks ``PATH_TO_URL_MAP`` looking for a matching prefix in the + *resolved* path string. Returns the original path unchanged if no + match is found. + """ + resolved = str(Path(nc_file_path).resolve()) + for local_prefix, url_prefix in PATH_TO_URL_MAP.items(): + if resolved.startswith(local_prefix): + return resolved.replace(local_prefix, url_prefix, 1) + return resolved + + +def get_git_url(script_name: str, version: str) -> str: + """Return a GitHub web URL for *script_name* at *version*.""" + return f"{GIT_WEB_BASE}/{version}/{script_name}" + + +def _find_first( + session: requests.Session, + endpoint: str, + params: dict, + api_base: str = SSDS_API_BASE, +) -> dict | None: + """GET /{endpoint}/ filtered by *params* — return first result or ``None``.""" + resp = session.get(f"{api_base}/{endpoint}/", params=params, timeout=REQUEST_TIMEOUT) + resp.raise_for_status() + data = resp.json() + # DRF pagination wraps list results in {"count": N, "results": [...]} + results = data.get("results", data) if isinstance(data, dict) else data + return results[0] if isinstance(results, list) and results else None + + +def find_person_by_email( + session: requests.Session, + email: str, + api_base: str = SSDS_API_BASE, +) -> dict | None: + """``GET /persons?email=…`` — return the first match or ``None``.""" + return _find_first(session, "persons", {"email": email}, api_base) + + +def find_software( + session: requests.Session, + name: str, + softwareversion: str, + uristring: str, + api_base: str = SSDS_API_BASE, +) -> dict | None: + "``GET /software?name=…&softwareversion=…&uristring=…`` — return the first match or ``None``." + return _find_first( + session, + "software", + {"name": name, "softwareversion": softwareversion, "uristring": uristring}, + api_base, + ) + + +def find_datacontainer( + session: requests.Session, + uristring: str, + api_base: str = SSDS_API_BASE, +) -> dict | None: + """``GET /datacontainers?uristring=…`` — return the first match or ``None``.""" + return _find_first(session, "datacontainers", {"uristring": uristring}, api_base) + + +def find_resource( + session: requests.Session, + uristring: str, + api_base: str = SSDS_API_BASE, +) -> dict | None: + """``GET /resources?uristring=…`` — return the first match or ``None``.""" + return _find_first(session, "resources", {"uristring": uristring}, api_base) + + +def _create_entity( + session: requests.Session, + endpoint: str, + payload: dict, + api_base: str = SSDS_API_BASE, +) -> dict: + """POST *payload* to *endpoint* and return the created entity.""" + url = f"{api_base}/{endpoint}/" + resp = session.post(url, json=payload, timeout=REQUEST_TIMEOUT) + if not resp.ok: + err_txt = ( + f"{resp.status_code} {resp.reason} for url: {url}" + f"\n payload: {payload}\n response: {resp.text}" + ) + raise requests.HTTPError(err_txt, response=resp) + return resp.json() + + +def _find_process_run_by_output( + session: requests.Session, + dods_url: str, + api_base: str = SSDS_API_BASE, +) -> dict | None: + """``GET /dataproducers?output_uri=…`` — find an existing ProcessRun.""" + return _find_first(session, "dataproducers", {"output_uri": dods_url}, api_base) + + +def _ensure_link( + session: requests.Session, + endpoint: str, + params: dict, + api_base: str = SSDS_API_BASE, +) -> None: + """POST the association; silently ignore duplicate-key responses. + + Prefer POST-first over a GET-check-then-POST pattern: junction-table + endpoints may not support filtering by FK query params, which would + cause _find_first to return a false positive and silently skip the POST. + """ + url = f"{api_base}/{endpoint}/" + resp = session.post(url, json=params, timeout=REQUEST_TIMEOUT) + if resp.ok: + return + # 409 Conflict, 400, or 500 whose body signals a unique/PK constraint violation. + # Django can surface IntegrityError as a 500 when the view lacks explicit handling. + _dup_markers = ("unique", "already exists", "duplicate key", "primary key constraint") + body = resp.text.lower() + if resp.status_code == requests.codes.conflict or any(m in body for m in _dup_markers): + return + err_txt = ( + f"{resp.status_code} {resp.reason} for url: {url}" + f"\n payload: {params}\n response: {resp.text}" + ) + raise requests.HTTPError(err_txt, response=resp) + + +# --------------------------------------------------------------------------- +# Core function +# --------------------------------------------------------------------------- +def submit_process_run( # noqa: PLR0913 C901 + nc_file_path: str, + input_uris: list[str], + *, + poc_email: str = "mccann@mbari.org", + pr_start: str | None = None, + pr_end: str | None = None, + software_name: str = "auv-python", + software_description: str = "AUV data processing pipeline", + software_version: str | None = None, + script_name: str = "src/data/process.py", + cmd_line_args: str = "", + log_file_url: str | None = None, + additional_resources: list[dict] | None = None, + api_key: str | None = None, + api_key_header: str | None = None, + api_base: str = SSDS_API_BASE, + session: requests.Session | None = None, + log: logging.Logger | None = None, +) -> dict | None: + """Submit a ProcessRun provenance record to the SSDS_Metadata database. + + This is the direct Python equivalent of the legacy Perl function + ``submitDStoNetCDFProcessRun``. + + Returns the ProcessRun entity dict on success, or ``None`` on failure. + """ + log = log or logger + session = session or build_authenticated_session( + api_key=api_key, + api_key_header=api_key_header, + ) + resolved_api_key_header = ( + api_key_header or os.environ.get(ENV_SSDS_API_KEY_HEADER) or "X-API-Key" + ) + if resolved_api_key_header not in session.headers: + log.warning( + "No SSDS API key header detected on session; set %s " + "(or pass a pre-authenticated session). " + "If the server uses cookie/SSO auth, this warning can be ignored.", + ENV_SSDS_API_KEY, + ) + software_version = software_version or _get_git_version() + + # a. Look up Person ------------------------------------------------------- + person = find_person_by_email(session, poc_email, api_base) + if person is None: + log.warning("Person with email %s not found – skipping provenance", poc_email) + return None + + # b. Build Resource list -------------------------------------------------- + resources: list[dict] = [] + # Script URL + script_url = get_git_url(script_name, software_version) + resources.append( + { + "name": "processing_script", + "uristring": script_url, + "description": f"Processing script: {script_name}", + } + ) + # Command line + if cmd_line_args: + resources.append( + { + "name": "command_line", + "uristring": f"urn:cmdline:{gethostname()}", + "description": cmd_line_args, + } + ) + # Processing log + if log_file_url: + resources.append( + { + "name": "processing_log", + "uristring": log_file_url, + "description": "Processing log file", + } + ) + if additional_resources: + resources.extend(additional_resources) + + # c. Create Software entity ----------------------------------------------- + software_payload = { + "name": software_name, + "version": "1", + "description": software_description, + "softwareversion": software_version, + "uristring": f"{GITHUB_REPO}/tree/{software_version}", + } + software = find_software( + session, software_name, software_version, software_payload["uristring"], api_base + ) or _create_entity(session, "software", software_payload, api_base) + + # d. Get or create ProcessRun (DataProducer) ------------------------------ + output_url = get_dods_url(nc_file_path) + existing = _find_process_run_by_output(session, output_url, api_base) + if existing: + pr_id = existing["id"] + # Update with new timing / description + update_payload = { + "description": f"Processing run for {Path(nc_file_path).name}", + "startdate": pr_start, + "enddate": pr_end, + "personid_fk": person["id"], + "softwareid_fk": software["id"], + } + url = f"{api_base}/dataproducers/{pr_id}/" + session.put(url, json=update_payload, timeout=REQUEST_TIMEOUT).raise_for_status() + process_run = {**existing, **update_payload} + log.info("Updated existing ProcessRun id=%s", pr_id) + else: + process_run = _create_entity( + session, + "dataproducers", + { + "name": f"Processing {Path(nc_file_path).name}", + "description": f"Processing run for {Path(nc_file_path).name}", + "version": "1", + "startdate": pr_start, + "enddate": pr_end, + "dataproducertype": "ProcessRun", + "personid_fk": person["id"], + "softwareid_fk": software["id"], + }, + api_base, + ) + pr_id = process_run["id"] + log.info("Created new ProcessRun id=%s", pr_id) + + # e. Get or create output DataContainer ------------------------------------ + output_dc = find_datacontainer(session, output_url, api_base) + if output_dc is None: + output_dc = _create_entity( + session, + "datacontainers", + { + "name": Path(nc_file_path).name, + "version": "1", + "uristring": output_url, + "description": f"Processed NetCDF output: {Path(nc_file_path).name}", + "dataproducerid_fk": pr_id, + }, + api_base, + ) + elif output_dc.get("dataproducerid_fk") != pr_id: + dc_id = output_dc["id"] + session.patch( + f"{api_base}/datacontainers/{dc_id}/", + json={"dataproducerid_fk": pr_id}, + timeout=REQUEST_TIMEOUT, + ).raise_for_status() + log.info("Linked existing DataContainer id=%s to ProcessRun id=%s", dc_id, pr_id) + + # f. Link relationships --------------------------------------------------- + + # Inputs — get or create DataContainer, then link via dataproducer-inputs + for uri in input_uris: + input_dc = find_datacontainer(session, uri, api_base) or _create_entity( + session, + "datacontainers", + { + "name": Path(uri).name, + "version": "1", + "uristring": uri, + "description": f"Input: {Path(uri).name}", + }, + api_base, + ) + _ensure_link( + session, + "dataproducer-inputs", + {"datacontainerid_fk": input_dc["id"], "dataproducerid_fk": pr_id}, + api_base, + ) + + # Resources — create each Resource then link via dataproducer-assoc-resource + for res_payload in resources: + res = find_resource(session, res_payload["uristring"], api_base) or _create_entity( + session, "resources", {**res_payload, "version": "1"}, api_base + ) + _ensure_link( + session, + "dataproducer-assoc-resource", + {"dataproducerid_fk": pr_id, "resourceid_fk": res["id"]}, + api_base, + ) + + # g. Log result ----------------------------------------------------------- + log.info("Provenance recorded: %s", f"{api_base}/dataproducers/{pr_id}") + return process_run + + +# --------------------------------------------------------------------------- +# CLI +# --------------------------------------------------------------------------- +def process_command_line() -> argparse.Namespace: + parser = argparse.ArgumentParser( + description=__doc__, + formatter_class=argparse.RawDescriptionHelpFormatter, + ) + parser.add_argument("nc_file", help="Path to the output NetCDF file") + parser.add_argument( + "--input", + dest="inputs", + action="append", + default=[], + help="Input URI (may be repeated)", + ) + parser.add_argument("--email", default="mccann@mbari.org", help="Point-of-contact email") + parser.add_argument("--script", default="src/data/process.py", help="Processing script path") + parser.add_argument("--cmd", default=" ".join(sys.argv), help="Command line string") + parser.add_argument("--log-url", default=None, help="URL to processing log") + parser.add_argument( + "--api-base", + default=SSDS_API_BASE, + help="SSDS REST API base URL", + ) + parser.add_argument("-v", "--verbose", action="count", default=0) + return parser.parse_args() + + +if __name__ == "__main__": + args = process_command_line() + logging.basicConfig( + level=logging.DEBUG if args.verbose else logging.INFO, + format="%(asctime)s %(levelname)s %(name)s: %(message)s", + ) + now = datetime.now(tz=UTC).isoformat() + result = submit_process_run( + nc_file_path=args.nc_file, + input_uris=args.inputs, + poc_email=args.email, + pr_start=now, + pr_end=now, + script_name=args.script, + cmd_line_args=args.cmd, + log_file_url=args.log_url, + api_base=args.api_base, + ) + sys.exit(0 if result else 1) From 0a06712990d788616c281c7cdb7a2abacfb2f638 Mon Sep 17 00:00:00 2001 From: Mike McCann Date: Thu, 2 Apr 2026 10:37:29 -0700 Subject: [PATCH 4/6] Simplify provenance.py using new /api/process-runs/ batch endpoint MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Replace the multi-step CRUD orchestration (10–15 round trips: person lookup, software find-or-create, ProcessRun get-or-create, DataContainer get-or-create, input/resource linking) with a single POST to the new /api/process-runs/ endpoint. The server now resolves or creates all related entities atomically from natural keys (person_email, software_name/version, output_uri, input_uris, resources). Remove all client-side helpers that are no longer needed: _find_first, find_person_by_email, find_software, find_datacontainer, find_resource, _create_entity, _find_process_run_by_output, _ensure_link. Net change: -214 lines. Public API (get_dods_url, submit_process_run) and all call sites in process.py are unchanged. Closes #144 (see also: previous junction-table linking bugs fixed while this endpoint was being developed on the server side) --- src/data/provenance.py | 276 +++++------------------------------------ 1 file changed, 31 insertions(+), 245 deletions(-) diff --git a/src/data/provenance.py b/src/data/provenance.py index 1a33bd2..9bcdb2f 100644 --- a/src/data/provenance.py +++ b/src/data/provenance.py @@ -1,8 +1,8 @@ """Record data processing provenance in the SSDS_Metadata database. -Replaces the legacy Perl ``submitDStoNetCDFProcessRun`` with authenticated -REST API calls to the mooring-ssds service. Every public helper is -designed to be called from ``process.py`` after a successful processing +Submits a single POST to ``/api/process-runs/``, which atomically resolves +or creates all related entities (Person, Software, DataContainers, Resources) +server-side. Can be called from ``process.py`` after a successful processing run, or standalone via the CLI at the bottom of this file. See: https://github.com/mbari-org/auv-python/issues/144 @@ -27,7 +27,6 @@ ENV_SSDS_API_BASE = "SSDS_API_BASE" SSDS_API_BASE = os.environ.get(ENV_SSDS_API_BASE, "https://mooring-ssds.shore.mbari.org/api") GIT_WEB_BASE = "https://github.com/mbari-org/auv-python/blob" -GITHUB_REPO = "https://github.com/mbari-org/auv-python" REQUEST_TIMEOUT = 30 # Environment variable names used for authentication stubs. @@ -103,124 +102,10 @@ def get_git_url(script_name: str, version: str) -> str: return f"{GIT_WEB_BASE}/{version}/{script_name}" -def _find_first( - session: requests.Session, - endpoint: str, - params: dict, - api_base: str = SSDS_API_BASE, -) -> dict | None: - """GET /{endpoint}/ filtered by *params* — return first result or ``None``.""" - resp = session.get(f"{api_base}/{endpoint}/", params=params, timeout=REQUEST_TIMEOUT) - resp.raise_for_status() - data = resp.json() - # DRF pagination wraps list results in {"count": N, "results": [...]} - results = data.get("results", data) if isinstance(data, dict) else data - return results[0] if isinstance(results, list) and results else None - - -def find_person_by_email( - session: requests.Session, - email: str, - api_base: str = SSDS_API_BASE, -) -> dict | None: - """``GET /persons?email=…`` — return the first match or ``None``.""" - return _find_first(session, "persons", {"email": email}, api_base) - - -def find_software( - session: requests.Session, - name: str, - softwareversion: str, - uristring: str, - api_base: str = SSDS_API_BASE, -) -> dict | None: - "``GET /software?name=…&softwareversion=…&uristring=…`` — return the first match or ``None``." - return _find_first( - session, - "software", - {"name": name, "softwareversion": softwareversion, "uristring": uristring}, - api_base, - ) - - -def find_datacontainer( - session: requests.Session, - uristring: str, - api_base: str = SSDS_API_BASE, -) -> dict | None: - """``GET /datacontainers?uristring=…`` — return the first match or ``None``.""" - return _find_first(session, "datacontainers", {"uristring": uristring}, api_base) - - -def find_resource( - session: requests.Session, - uristring: str, - api_base: str = SSDS_API_BASE, -) -> dict | None: - """``GET /resources?uristring=…`` — return the first match or ``None``.""" - return _find_first(session, "resources", {"uristring": uristring}, api_base) - - -def _create_entity( - session: requests.Session, - endpoint: str, - payload: dict, - api_base: str = SSDS_API_BASE, -) -> dict: - """POST *payload* to *endpoint* and return the created entity.""" - url = f"{api_base}/{endpoint}/" - resp = session.post(url, json=payload, timeout=REQUEST_TIMEOUT) - if not resp.ok: - err_txt = ( - f"{resp.status_code} {resp.reason} for url: {url}" - f"\n payload: {payload}\n response: {resp.text}" - ) - raise requests.HTTPError(err_txt, response=resp) - return resp.json() - - -def _find_process_run_by_output( - session: requests.Session, - dods_url: str, - api_base: str = SSDS_API_BASE, -) -> dict | None: - """``GET /dataproducers?output_uri=…`` — find an existing ProcessRun.""" - return _find_first(session, "dataproducers", {"output_uri": dods_url}, api_base) - - -def _ensure_link( - session: requests.Session, - endpoint: str, - params: dict, - api_base: str = SSDS_API_BASE, -) -> None: - """POST the association; silently ignore duplicate-key responses. - - Prefer POST-first over a GET-check-then-POST pattern: junction-table - endpoints may not support filtering by FK query params, which would - cause _find_first to return a false positive and silently skip the POST. - """ - url = f"{api_base}/{endpoint}/" - resp = session.post(url, json=params, timeout=REQUEST_TIMEOUT) - if resp.ok: - return - # 409 Conflict, 400, or 500 whose body signals a unique/PK constraint violation. - # Django can surface IntegrityError as a 500 when the view lacks explicit handling. - _dup_markers = ("unique", "already exists", "duplicate key", "primary key constraint") - body = resp.text.lower() - if resp.status_code == requests.codes.conflict or any(m in body for m in _dup_markers): - return - err_txt = ( - f"{resp.status_code} {resp.reason} for url: {url}" - f"\n payload: {params}\n response: {resp.text}" - ) - raise requests.HTTPError(err_txt, response=resp) - - # --------------------------------------------------------------------------- # Core function # --------------------------------------------------------------------------- -def submit_process_run( # noqa: PLR0913 C901 +def submit_process_run( # noqa: PLR0913 nc_file_path: str, input_uris: list[str], *, @@ -228,7 +113,6 @@ def submit_process_run( # noqa: PLR0913 C901 pr_start: str | None = None, pr_end: str | None = None, software_name: str = "auv-python", - software_description: str = "AUV data processing pipeline", software_version: str | None = None, script_name: str = "src/data/process.py", cmd_line_args: str = "", @@ -240,12 +124,12 @@ def submit_process_run( # noqa: PLR0913 C901 session: requests.Session | None = None, log: logging.Logger | None = None, ) -> dict | None: - """Submit a ProcessRun provenance record to the SSDS_Metadata database. + """Submit a ProcessRun to the SSDS_Metadata database via POST /api/process-runs/. - This is the direct Python equivalent of the legacy Perl function - ``submitDStoNetCDFProcessRun``. + The server resolves or creates all related entities (Person, Software, + DataContainers, Resources) from the supplied natural keys. - Returns the ProcessRun entity dict on success, or ``None`` on failure. + Returns the created ProcessRun dict on success, or raises on HTTP error. """ log = log or logger session = session or build_authenticated_session( @@ -264,24 +148,13 @@ def submit_process_run( # noqa: PLR0913 C901 ) software_version = software_version or _get_git_version() - # a. Look up Person ------------------------------------------------------- - person = find_person_by_email(session, poc_email, api_base) - if person is None: - log.warning("Person with email %s not found – skipping provenance", poc_email) - return None - - # b. Build Resource list -------------------------------------------------- - resources: list[dict] = [] - # Script URL - script_url = get_git_url(script_name, software_version) - resources.append( + resources: list[dict] = [ { "name": "processing_script", - "uristring": script_url, + "uristring": get_git_url(script_name, software_version), "description": f"Processing script: {script_name}", - } - ) - # Command line + }, + ] if cmd_line_args: resources.append( { @@ -290,7 +163,6 @@ def submit_process_run( # noqa: PLR0913 C901 "description": cmd_line_args, } ) - # Processing log if log_file_url: resources.append( { @@ -302,114 +174,28 @@ def submit_process_run( # noqa: PLR0913 C901 if additional_resources: resources.extend(additional_resources) - # c. Create Software entity ----------------------------------------------- - software_payload = { - "name": software_name, - "version": "1", - "description": software_description, - "softwareversion": software_version, - "uristring": f"{GITHUB_REPO}/tree/{software_version}", + payload = { + "output_uri": get_dods_url(nc_file_path), + "output_name": Path(nc_file_path).name, + "input_uris": input_uris, + "software_name": software_name, + "software_version": software_version, + "person_email": poc_email, + "startdate": pr_start, + "enddate": pr_end, + "resources": resources, } - software = find_software( - session, software_name, software_version, software_payload["uristring"], api_base - ) or _create_entity(session, "software", software_payload, api_base) - - # d. Get or create ProcessRun (DataProducer) ------------------------------ - output_url = get_dods_url(nc_file_path) - existing = _find_process_run_by_output(session, output_url, api_base) - if existing: - pr_id = existing["id"] - # Update with new timing / description - update_payload = { - "description": f"Processing run for {Path(nc_file_path).name}", - "startdate": pr_start, - "enddate": pr_end, - "personid_fk": person["id"], - "softwareid_fk": software["id"], - } - url = f"{api_base}/dataproducers/{pr_id}/" - session.put(url, json=update_payload, timeout=REQUEST_TIMEOUT).raise_for_status() - process_run = {**existing, **update_payload} - log.info("Updated existing ProcessRun id=%s", pr_id) - else: - process_run = _create_entity( - session, - "dataproducers", - { - "name": f"Processing {Path(nc_file_path).name}", - "description": f"Processing run for {Path(nc_file_path).name}", - "version": "1", - "startdate": pr_start, - "enddate": pr_end, - "dataproducertype": "ProcessRun", - "personid_fk": person["id"], - "softwareid_fk": software["id"], - }, - api_base, - ) - pr_id = process_run["id"] - log.info("Created new ProcessRun id=%s", pr_id) - - # e. Get or create output DataContainer ------------------------------------ - output_dc = find_datacontainer(session, output_url, api_base) - if output_dc is None: - output_dc = _create_entity( - session, - "datacontainers", - { - "name": Path(nc_file_path).name, - "version": "1", - "uristring": output_url, - "description": f"Processed NetCDF output: {Path(nc_file_path).name}", - "dataproducerid_fk": pr_id, - }, - api_base, - ) - elif output_dc.get("dataproducerid_fk") != pr_id: - dc_id = output_dc["id"] - session.patch( - f"{api_base}/datacontainers/{dc_id}/", - json={"dataproducerid_fk": pr_id}, - timeout=REQUEST_TIMEOUT, - ).raise_for_status() - log.info("Linked existing DataContainer id=%s to ProcessRun id=%s", dc_id, pr_id) - - # f. Link relationships --------------------------------------------------- - - # Inputs — get or create DataContainer, then link via dataproducer-inputs - for uri in input_uris: - input_dc = find_datacontainer(session, uri, api_base) or _create_entity( - session, - "datacontainers", - { - "name": Path(uri).name, - "version": "1", - "uristring": uri, - "description": f"Input: {Path(uri).name}", - }, - api_base, - ) - _ensure_link( - session, - "dataproducer-inputs", - {"datacontainerid_fk": input_dc["id"], "dataproducerid_fk": pr_id}, - api_base, - ) - # Resources — create each Resource then link via dataproducer-assoc-resource - for res_payload in resources: - res = find_resource(session, res_payload["uristring"], api_base) or _create_entity( - session, "resources", {**res_payload, "version": "1"}, api_base - ) - _ensure_link( - session, - "dataproducer-assoc-resource", - {"dataproducerid_fk": pr_id, "resourceid_fk": res["id"]}, - api_base, + url = f"{api_base}/process-runs/" + resp = session.post(url, json=payload, timeout=REQUEST_TIMEOUT) + if not resp.ok: + err_txt = ( + f"{resp.status_code} {resp.reason} for url: {url}" + f"\n payload: {payload}\n response: {resp.text}" ) - - # g. Log result ----------------------------------------------------------- - log.info("Provenance recorded: %s", f"{api_base}/dataproducers/{pr_id}") + raise requests.HTTPError(err_txt, response=resp) + process_run = resp.json() + log.info("Provenance recorded: %s -> id=%s", url, process_run.get("id", "?")) return process_run From 5d3688b6ca39584d8f0f92cfb1db6044338d5a85 Mon Sep 17 00:00:00 2001 From: Mike McCann Date: Thu, 2 Apr 2026 10:52:22 -0700 Subject: [PATCH 5/6] Update EXPECTED_SIZE_GITHUB --- src/data/test_process_i2map.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/data/test_process_i2map.py b/src/data/test_process_i2map.py index 88d0026..89eb695 100644 --- a/src/data/test_process_i2map.py +++ b/src/data/test_process_i2map.py @@ -30,7 +30,7 @@ def test_process_i2map(complete_i2map_processing): # but it will alert us if a code change unexpectedly changes the file size. # If code changes are expected to change the file size then we should # update the expected size here. - EXPECTED_SIZE_GITHUB = 63137 + EXPECTED_SIZE_GITHUB = 63130 EXPECTED_SIZE_ACT = 63106 EXPECTED_SIZE_LOCAL = 64650 if str(proc.args.base_path).startswith("/home/runner"): From d994d10a4cb1276b065a3b78a59a76fa8dda4c4e Mon Sep 17 00:00:00 2001 From: Mike McCann Date: Fri, 3 Apr 2026 15:35:51 -0700 Subject: [PATCH 6/6] Add producer, software URL, and OPeNDAP links to provenance payload - Add producer_name and producer_description to submit_process_run() and populate them with execution context in _submit_provenance() - Set software_uristring to the GitHub tree URL at the commit hash - Add output_dodsurlstring and input_dodsurlstrings with .html suffix for browsable OPeNDAP URLs - Refactor _submit_provenance() to accept base_path separately from output_nc, building the full path internally to match --log_file convention and produce correct OPeNDAP URLs --- src/data/process.py | 41 ++++++++++++++++++++++++++++------------- src/data/provenance.py | 11 +++++++++-- 2 files changed, 37 insertions(+), 15 deletions(-) diff --git a/src/data/process.py b/src/data/process.py index e8b6e64..33e881e 100755 --- a/src/data/process.py +++ b/src/data/process.py @@ -713,20 +713,33 @@ def create_products(self, mission: str = None, log_file: str = None) -> None: def _submit_provenance( # noqa: PLR0913 self, output_nc: str, + base_path: str, input_files: list[str], pr_start: str, pr_end: str, script_name: str = "src/data/process.py", log_file: str | None = None, ) -> None: - """Submit a provenance record — failures are logged, never raised.""" + """Submit a provenance record — failures are logged, never raised. + + *output_nc* is a path relative to *base_path* (e.g. + ``ahi/missionlogs/.../file_1S.nc``), matching the ``--log_file`` + convention. *base_path* is used internally to locate the file and + build the OPeNDAP URL. + """ try: - if not Path(output_nc).exists(): - self.logger.debug("Output %s not found, skipping provenance", output_nc) + full_nc = str(Path(base_path, output_nc)) + if not Path(full_nc).exists(): + self.logger.debug("Output %s not found, skipping provenance", full_nc) return log_url = get_dods_url(log_file) if log_file else None submit_process_run( - nc_file_path=output_nc, + producer_name=( + f"auv-python - Execution of {Path(script_name).name}" + f" to produce {Path(output_nc)}" + ), + producer_description=self.commandline, + nc_file_path=full_nc, input_uris=input_files, pr_start=pr_start, pr_end=pr_end, @@ -886,15 +899,16 @@ def process_mission(self, mission: str, src_dir: str = "") -> None: # noqa: C90 self.resample(mission) self.create_products(mission) if self.config["update_ssds_provenance"]: - resampled = Path( - self.config["base_path"], - self.auv_name, - MISSIONNETCDFS, - mission, - f"{self.auv_name}_{mission}_{FREQ}.nc", - ) self._submit_provenance( - output_nc=str(resampled), + output_nc=str( + Path( + self.auv_name, + MISSIONNETCDFS, + mission, + f"{self.auv_name}_{mission}_{FREQ}.nc", + ) + ), + base_path=self.config["base_path"], input_files=[ str( Path( @@ -1141,7 +1155,8 @@ def process_log_file(self, log_file: str) -> None: self.create_products(log_file=log_file) if self.config["update_ssds_provenance"]: self._submit_provenance( - output_nc=str(resampled_file), + output_nc=str(Path(log_file).parent / f"{Path(log_file).stem}_{FREQ}.nc"), + base_path=str(BASE_LRAUV_PATH), input_files=[get_dods_url(str(Path(BASE_LRAUV_PATH, log_file)))], pr_start=_pr_start, pr_end=datetime.now(tz=UTC).isoformat(), diff --git a/src/data/provenance.py b/src/data/provenance.py index 9bcdb2f..fd8df76 100644 --- a/src/data/provenance.py +++ b/src/data/provenance.py @@ -109,6 +109,8 @@ def submit_process_run( # noqa: PLR0913 nc_file_path: str, input_uris: list[str], *, + producer_name: str | None = None, + producer_description: str | None = None, poc_email: str = "mccann@mbari.org", pr_start: str | None = None, pr_end: str | None = None, @@ -174,12 +176,17 @@ def submit_process_run( # noqa: PLR0913 if additional_resources: resources.extend(additional_resources) + output_uri = get_dods_url(nc_file_path) payload = { - "output_uri": get_dods_url(nc_file_path), - "output_name": Path(nc_file_path).name, + "output_uri": output_uri, + "output_dodsurlstring": f"{output_uri}.html", + "producer_name": producer_name, + "producer_description": producer_description, "input_uris": input_uris, + "input_dodsurlstrings": [f"{uri}.html" for uri in input_uris], "software_name": software_name, "software_version": software_version, + "software_uristring": f"https://github.com/mbari-org/auv-python/tree/{software_version}", "person_email": poc_email, "startdate": pr_start, "enddate": pr_end,