diff --git a/README.md b/README.md index bba79ff..d9ed740 100644 --- a/README.md +++ b/README.md @@ -174,6 +174,57 @@ exact same arguments. - `catalog.close()` (or `with Catalog(...) as catalog:`) — disposes the REST/Flight session(s). Idempotent, and covered by a process-exit/SIGTERM fallback if you forget. +## SDI metadata + +`eea_datalakehouse.sdi` copies a dataset's ISO 19115-3 metadata record from the EEA SDI +catalogue (GeoNetwork) into the dataset's `metadata` folder in DDS. The XML is a **document +upload**, not an ingest: nothing becomes a Dremio table. The calling code always passes the +SDI UUID. See `docs/sdi-integration-plan.md`. + +```python +from eea_datalakehouse.sdi import SdiController + +with SdiController.from_env() as sdi: + metadata = sdi.get_xml("070d9baa-448d-4168-8514-7dadb3ad876d") # SDI -> bytes + result = sdi.push_to_dds(metadata, "catalog/water_management_resources/bathing_water/bwd") + result.dds_path # ".../bwd/metadata/070d9baa-448d-4168-8514-7dadb3ad876d.xml" + result.action # "uploaded" | "unchanged" | "replaced" +``` + +- `get_xml(uuid)` checks the response is `mdb:MD_Metadata` and identifies `uuid` + (`NotIso19115_3`, `UuidMismatch`; `SdiNotFound` / `SdiAuthError` for 404 / 401-403). +- `push_to_dds(metadata, dds_path, folder="metadata", force=False)` compares with the copy + already in DDS: same bytes is `unchanged`; an older copy is `replaced`; a copy that is newer or + undated (someone edited it in DDS) raises `DdsCopyConflict` unless `force=True`. +- `resolve_series(series_uuid)` returns the one release of a series that is not superseded, or + raises `NotCurrentError` listing every release. +- `metadata.provenance_tags()` gives `sdi_record_uuid` / `sdi_edition` / `sdi_date_stamp` in + `Catalog.setmeta2wiki`'s `tags=` shape, if you want them on the dataset's wiki. + +Environment: `SDI_API_URL` (empty means `https://sdi.eea.europa.eu/catalogue`), optional +`SDI_USERNAME` / `SDI_PASSWORD` for non-public records, and for the upload `DDS_BASE_URL` plus +`_DREMIO_USER` / `_DREMIO_PWD`, as for ingest. The DDS document endpoints +(`dds_documents/client.py`) are still assumptions to confirm with the DDS team. + +In a notebook, two magics split the work: `%sdi` (`SdiSession`) reads SDI, and `%metadata` +(`MetadataSession`) pushes metadata files to DDS. Both are built on first use from the kernel +environment like `%catalog`'s and `%ingest`'s sessions. `%metadata push_to_dds` uploads the +record `%sdi get_xml` fetched last, so it needs only the path (worked example: +`debugger/sdi_session_example.ipynb`): + +```python +import eea_datalakehouse.notebook # registers %catalog/%ingest/%sdi/%metadata + +%sdi help +%sdi get_xml("070d9baa-448d-4168-8514-7dadb3ad876d") +%metadata push_to_dds("catalog/water_management_resources/bathing_water/bwd") +``` + +For the push, `DDS_BASE_URL` is read from the first `.env` in the notebook's folder or a parent +when `%metadata` builds its session (the kernel environment's value only applies if no `.env` +sets it; `%metadata dds_base_url()` shows which is in use). The Dremio identity is `_DREMIO_USER` / +`_DREMIO_PWD` (as `%ingest`), falling back to `DREMIO_USERNAME` / `DREMIO_TOKEN` (as `%catalog`). + ## Notebook facade (`%catalog` / `%ingest`) For interactive use in JupyterLab, `eea_datalakehouse.notebook` registers two line magics @@ -281,6 +332,11 @@ context itself — raises `CatalogSessionError` if none is set yet: | `dds_ingestion/client.py` | thin, unit-testable HTTP client (`IngestClient`) | | `dds_ingestion/progress.py` | tqdm progress bar with graceful fallback | | `dds_ingestion/folder.py` | `FolderIngest` orchestration (scan/parallel/resume) | +| `dds_documents/client.py` | `DocumentsClient` — put/get/list plain files in DDS folders (no ingest) | +| `sdi/controller.py` | `SdiController` — `get_xml()`, `push_to_dds()`, `resolve_series()` | +| `sdi/catalogue.py` | `SdiCatalogue` — read-only GeoNetwork REST client | +| `sdi/session.py` | `SdiSession` / `MetadataSession` — what `%sdi` / `%metadata` dispatch onto | +| `sdi/iso.py` | reads UUID, title, edition, dates and series children from ISO 19115-3 | | `catalog/client.py` | `Catalog` — two connections (REST + Flight), every operation as a method | | `catalog/operations.py` | the operations themselves (table2view, datacopy, createfolder, ...), as functions taking an executor and/or a `CatalogRestClient` | | `catalog/sql.py` | `SqlExecutor` protocol + REST/Flight implementations, `resolve_executor()` (transport env var) | @@ -360,7 +416,7 @@ To pin to one specific release instead, use the exact tag the live badges under [Releasing a new version](#releasing-a-new-version)): ```python -%pip install "git+https://github.com/eeadata/EEALakeHouse.python.git@v0.1.17" +%pip install "git+https://github.com/eeadata/EEALakeHouse.python.git@v0.1.19" # staging's latest release (early access) %pip install "git+https://github.com/eeadata/EEALakeHouse.python.git@v0.1.19-staging" diff --git a/debugger/.env.example b/debugger/.env.example index 386742b..c5d034c 100644 --- a/debugger/.env.example +++ b/debugger/.env.example @@ -11,10 +11,21 @@ SDI_METADATA_DIR= DDS_BASE_URL= DREMIO_BASE_URL= +# --- SDI catalogue (eea_datalakehouse.sdi) ------------------------------------ +# GeoNetwork base URL; the client appends /srv/api/... Empty means the public +# EEA catalogue, https://sdi.eea.europa.eu/catalogue. +SDI_API_URL= +# Optional: an SDI account, only needed to read records that are not public. +# Public records (all published EEA datasets) need neither. +SDI_USERNAME= +SDI_PASSWORD= + # --- credentials — never commit ------------------------------------------- # DREMIO_TOKEN is a Personal Access Token (Dremio → Account Settings → Personal # Access Tokens). In the local stack's stub auth mode the token is simply the # username: alice (read-write), bob (read-only), carol (no access). +# debug_run.py copies these to _DREMIO_USER / _DREMIO_PWD, which the DDS +# ingest and documents clients read. DREMIO_USER= DREMIO_TOKEN= # Your own Redmine (taskman) API key: taskman > My account > API access key @@ -22,6 +33,15 @@ DREMIO_TOKEN= # referenced as "refs #1234" / "fixes #1234". Never leaves your machine. REDMINE_API_KEY= +# --- Entra ID → Dremio token exchange (debug_run.py's PAT helper) ------------- +# DREMIO_HOST is the Dremio host name only, no scheme. ENTRA_CLIENT_SECRET is +# only needed for the service-principal mode (--mode sp). +DREMIO_HOST= +ENTRA_TENANT_ID= +ENTRA_CLIENT_ID= +ENTRA_CLIENT_SECRET= +ENTRA_SCOPE= + # --- notebooks -------------------------------------------------------------- JUPYTER_HOST_PORT= # Leave empty for a token-free lab on localhost; set one if the port is exposed. diff --git a/debugger/debug_run.py b/debugger/debug_run.py index 873a827..1726386 100644 --- a/debugger/debug_run.py +++ b/debugger/debug_run.py @@ -40,7 +40,9 @@ from pathlib import Path from xmlrpc.client import Boolean from eea_datalakehouse.catalog import Catalog +from eea_datalakehouse.dds_documents import DocumentsClient from eea_datalakehouse.dds_ingestion import DremioCreds, FolderIngest, IngestClient +from eea_datalakehouse.sdi import SdiCatalogue, SdiController, metadata_path from eea_datalakehouse.dds_ingestion.common.dremio_identity import (endpoint, dds_credentials, resolve as _resolve_identity) @@ -420,6 +422,24 @@ def run_catalog_bulk_close() -> None: print(f"catalog_bulk_close after a.closed={a._closed} b.closed={b._closed}") +def run_sdi_extract() -> None: + """Download the release's ISO 19115-3 XML from SDI and push it to DDS's + metadata folder. Extraction runs for real even in DRY_RUN (it is a public, + read-only SDI call); only the DDS upload is skipped.""" + uuid = CONFIG["sdi"]["latest_record_uuid"] + documents = DocumentsClient(DDS_BASE_URL, dds_credentials(DREMIO_USERNAME, DREMIO_TOKEN)) + with SdiController(SdiCatalogue.from_env(), documents) as sdi, documents: + metadata = sdi.get_xml(uuid) + print(f"sdi_extract {metadata}") + print(f" revised {metadata.record.date_stamp}") + if DRY_RUN: + print("DRY RUN — not uploading to DDS. Set DRY_RUN = False to run this for real.") + print(f" would put {metadata_path(SDI_DDS_PATH, uuid)}") + return + result = sdi.push_to_dds(metadata, SDI_DDS_PATH) + print(f"sdi_extract {result.action} {result.dds_path}") + + # -------------------------------------------------------------------------- # Step 1 — get an Entra ID JWT @@ -599,6 +619,12 @@ def create_pat(host: str, access_token: str, username: str, label: str, ) TABLE2VIEW_IDEMPOTENCY_KEY = "debug-table2view-testview" + # --- sdi_extract ------------------------------------------------------------ + # Where the dataset lives in DDS; push_to_dds adds /metadata/{uuid}.xml. + # Provisional until the DDS document path is agreed (docs/sdi-integration-plan.md). + _t = CONFIG["target"] + SDI_DDS_PATH = f"catalog/{_t['domain']}/{_t['subdomain']}/{_t['dataflow']}" + # --- services -------------------------------------------------------------- # Resolved most-authoritative-first by endpoint(): # 1. the JupyterLab "Dremio Catalog" settings panel — ddsServerUrl / dremioUrl. @@ -657,6 +683,8 @@ def create_pat(host: str, access_token: str, username: str, label: str, #run_createfolder() #run_deletefolder() + #run_sdi_extract() + diff --git a/debugger/sdi_session_example.ipynb b/debugger/sdi_session_example.ipynb new file mode 100644 index 0000000..77600cd --- /dev/null +++ b/debugger/sdi_session_example.ipynb @@ -0,0 +1,179 @@ +{ + "cells": [ + { + "cell_type": "markdown", + "metadata": {}, + "source": "# `%sdi` + `%metadata` — SDI metadata into DDS\n\nA test notebook for the `%sdi` and `%metadata` magics from `eea_datalakehouse.notebook.magics`.\n`%sdi` downloads a dataset's ISO 19115-3 metadata record from the EEA SDI catalogue;\n`%metadata` uploads it to the dataset's `metadata` folder in DDS — a **document upload**, never an ingest: nothing becomes a\nDremio table. See `docs/sdi-integration-plan.md` for the design.\n\n**Kernel environment** — read the same way `%catalog` and `%ingest` read theirs, when each magic's\nfirst call builds its session:\n\n| Variable | Needed for | Notes |\n|---|---|---|\n| `SDI_API_URL` | `%sdi get_xml`, `%sdi resolve_series` | empty or unset → `https://sdi.eea.europa.eu/catalogue` |\n| `SDI_USERNAME` / `SDI_PASSWORD` | non-public records only | optional |\n| `DDS_BASE_URL` | `%metadata push_to_dds` | read from the first `.env` in the notebook's folder or a parent (here `debugger/.env`); else the kernel environment |\n| `_DREMIO_USER` / `_DREMIO_PWD` | `%metadata push_to_dds` | as for `%ingest`; else `DREMIO_USERNAME` / `DREMIO_TOKEN`, as for `%catalog` |\n\n`%sdi` makes public, read-only calls and needs none of the DDS settings. Install the extra this\nneeds once: `pip install \"EEADataLakehouse[notebook]\"`.", + "id": "bc3d7aa0" + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": "import eea_datalakehouse.notebook # registers %catalog/%ingest/%sdi — no %load_ext needed", + "id": "201451f8" + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": "## Quick reference\n\n`%sdi help` and `%metadata help` render every command as an HTML table, the same format as\n`%catalog help` and `%ingest help`. They work before any environment variable is set.", + "id": "4db70eeb" + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": "%sdi help", + "id": "655f890a" + }, + { + "cell_type": "code", + "id": "1cb0d697", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": "%metadata help" + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": "## Settings for this test\n\nThe bathing water release used throughout the library's tests. The DDS path is\n**provisional** until the DDS document path convention is agreed; `%metadata push_to_dds` appends\n`/metadata/.xml` to it.", + "id": "200b71fe" + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": "RELEASE_UUID = \"070d9baa-448d-4168-8514-7dadb3ad876d\" # Bathing Water Directive, 2025 v1.0\nSERIES_UUID = \"c3858959-90da-4c1b-b9ca-492db0e514df\" # its series\nDDS_PATH = \"catalog/water_management_resources/bathing_water/bwd\"", + "id": "fb073e8a" + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": "## 1. Get the XML from SDI\n\nRuns immediately. The record is checked before it is returned: it must be ISO 19115-3\n(`mdb:MD_Metadata`) and its identifier must be `RELEASE_UUID`. Any failure prints one\n`sdi error: ...` line instead of a traceback.", + "id": "795e9151" + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": "metadata = %sdi get_xml(RELEASE_UUID)\nmetadata", + "id": "6269d370" + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": "What was parsed out of it — title, edition, scope and the metadata dates:", + "id": "ddf2bdb3" + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": "metadata.record", + "id": "230bfbc2" + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": "print(metadata.xml[:600].decode())", + "id": "9fce9388" + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": "## 2. A series instead of a release\n\n`resolve_series` returns the one release of a series that is not superseded. If there are zero\nor several it refuses to guess and lists them all. The bathing water series currently has\n**two** non-superseded releases, so expect an `sdi error` listing both — the cell after it\nuses the release UUID directly, which is the normal way.", + "id": "4d67a127" + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": "%sdi resolve_series(SERIES_UUID)", + "id": "647cddc0" + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": "## 3. Push it to DDS\n\n**Writes to DDS.** `%metadata push_to_dds` uploads the record `%sdi get_xml` fetched last, so\nonly the DDS path is needed. It compares with the copy already there first:\n\n| DDS already has | Result |\n|---|---|\n| nothing | `uploaded` |\n| the same bytes | `unchanged` (no upload) |\n| an older copy | `replaced` |\n| a newer, same-dated or non-ISO copy (edited in DDS) | refused, unless `force=True` |\n\nThe DDS document endpoints are still interface assumptions to confirm with the DDS team, so\nrun this against a test instance first. Set `RUN_PUSH = True` to run it.", + "id": "873f813c" + }, + { + "cell_type": "markdown", + "id": "cbb0491a", + "metadata": {}, + "source": "Check which DDS the push will go to first. `DDS_BASE_URL` is taken from the first `.env`\nfound in this notebook's folder or a parent (read once, when `%metadata` builds its session), and\nfrom the kernel environment only if no `.env` sets it:" + }, + { + "cell_type": "code", + "id": "416cd8ae", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": "%metadata dds_base_url()" + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": "RUN_PUSH = False\n\nif RUN_PUSH:\n result = %metadata push_to_dds(DDS_PATH)\n print(result)\nelse:\n from eea_datalakehouse.sdi import metadata_path\n print(\"RUN_PUSH is False — would upload to\", metadata_path(DDS_PATH, metadata.uuid))", + "id": "b962036b" + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": "Running the push a second time should report `unchanged`. To replace a copy that was edited in\nDDS, pass `force=True`; to push a record other than the last one, pass it explicitly:\n\n```\n%metadata push_to_dds(DDS_PATH, force=True)\n%metadata push_to_dds(DDS_PATH, metadata=other_metadata)\n```", + "id": "e4248c6b" + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": "## 4. Optional: record the provenance on the wiki\n\n`provenance_tags()` gives `sdi_record_uuid` / `sdi_edition` / `sdi_date_stamp` in the shape\n`Catalog.setmeta2wiki`'s `tags` take (a folder's wiki \"# Meta Data\" section). Writing them is\nup to you — neither magic writes to Dremio itself:\n\n```python\nfrom eea_datalakehouse.catalog import Catalog\n\nwith Catalog(DREMIO_BASE_URL, DREMIO_TOKEN, username=DREMIO_USERNAME) as catalog:\n catalog.setmeta2wiki(\n \"catalog.water_management_resources.bathing_water.bwd\",\n tags=metadata.provenance_tags(),\n overwrite=True, # merge with the tags already there\n idempotency_key=f\"sdi-provenance-{metadata.uuid}\",\n )\n```", + "id": "48ed678c" + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": "metadata.provenance_tags()", + "id": "16fa54c8" + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": "## What errors look like\n\n```\n%sdi get_xml(\"00000000-0000-0000-0000-000000000000\")\n# -> sdi error: SDI API error 404 on GET /srv/api/records/.../formatters/xml: no SDI record with UUID '...'\n\n%metadata push_to_dds(DDS_PATH) # in a kernel with no Dremio identity\n# -> metadata error: pushing to DDS needs a Dremio identity in the kernel environment: ...\n```", + "id": "217752bd" + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": "%sdi get_xml(\"00000000-0000-0000-0000-000000000000\")", + "id": "0f8edce6" + } + ], + "metadata": { + "kernelspec": { + "display_name": "Python 3", + "language": "python", + "name": "python3" + }, + "language_info": { + "name": "python", + "version": "3.11" + } + }, + "nbformat": 4, + "nbformat_minor": 5 +} diff --git a/debugger/test_catalog_magic.ipynb b/debugger/test_catalog_magic.ipynb index daed23a..f01f907 100644 --- a/debugger/test_catalog_magic.ipynb +++ b/debugger/test_catalog_magic.ipynb @@ -499,11 +499,29 @@ "source": [ "%catalog schema('bwd.draft.bw_assessment.assessments') # should be empty now, after the delete" ] + }, + { + "cell_type": "code", + "execution_count": 1, + "id": "ca2ba7e1", + "metadata": {}, + "outputs": [ + { + "name": "stdout", + "output_type": "stream", + "text": [ + "UUUU\n" + ] + } + ], + "source": [ + "print(\"UUUU\")" + ] } ], "metadata": { "kernelspec": { - "display_name": ".venv (3.12.3.final.0)", + "display_name": ".venv (3.12.3)", "language": "python", "name": "python3" }, diff --git a/docs/sdi-integration-plan.md b/docs/sdi-integration-plan.md new file mode 100644 index 0000000..0cd3f3a --- /dev/null +++ b/docs/sdi-integration-plan.md @@ -0,0 +1,385 @@ +# SDI integration — development plan — **PROPOSED** + +> Status: proposed, 2026-09-30. **Steps 1–3 built 2026-10-01** (`sdi/`, +> `dds_documents/`, with tests); step 0 (confirming the DDS document endpoints) +> and step 4 (a run against a DDS test instance) are still open. +> +> Scope (2026-10-01): **extraction from SDI and upload to DDS only.** +> Publication to SDI is deferred; see [Not in this plan](#not-in-this-plan). + +"SDI" here means the EEA Spatial Data Infrastructure catalogue: the "EEA +geospatial data catalogue", GeoNetwork **4.4.9**, at +`https://sdi.eea.europa.eu/catalogue`. It holds the ISO 19115-3 metadata record +for each dataset. + +## The short version + +Add an **SDI controller** to `eea_datalakehouse` with one route, **Extract** +(SDI → DDS): + +- It takes a dataset's SDI UUID from the calling code. +- It downloads that record's ISO 19115-3 XML from the SDI REST API. +- It uploads the XML to the dataset's `metadata/` folder in DDS. + +The SDI REST API base URL is read from `SDI_API_URL`, e.g. +`SDI_API_URL=https://sdi.eea.europa.eu/catalogue`. The client appends +`/srv/api/...`. + +The controller sits on three pieces: + +- **`sdi/catalogue.py`**: a read-only SDI REST client. +- **`dds_documents/`**: a DDS documents client, which this library doesn't + have yet. The existing DDS client, `dds_ingestion`, only covers + `/api/v1/ingest`. +- **`sdi/iso.py`**: a small ISO 19115-3 reader for the few fields the + controller checks (UUID, title, edition, dates). + +## Where things stand in this repo + +- **Before this plan, no SDI code.** `src/eea_datalakehouse/` had `catalog`, + `dds_ingestion` and `notebook`, and nothing talked to SDI. +- **Unused dependencies.** `pyproject.toml` carries `lxml` ("ISO 19115-3 + metadata records"), `beautifulsoup4` ("SDI 'direct download' is an HTML + landing page") and `xlrd` ("some older SDI deliveries"). Nothing imports + any of them. +- **The DDS connection already exists.** `dds_ingestion/credentials.py` + resolves `DDS_BASE_URL` and the Dremio credentials (`_DREMIO_USER` / + `_DREMIO_PWD`). `dds_ingestion/client.py` sends the PAT as + `Authorization: Bearer`. The documents client reuses both. +- **Wiki metadata can already be written.** `catalog.setmeta2wiki()`, + `getmetafromwiki()` and `settagsto()` are there to record provenance on a + Dremio entity. No SDI data feeds them yet. +- **No SDI configuration yet.** `debugger/.env.example` has no SDI entries. + +## Verified against the live SDI catalogue (2026-09-30) + +These are public, read-only calls, using the bathing water series `c3858959…` +and its release `070d9baa…` as the example. + +- **`GET {SDI_API_URL}/srv/api/records/{uuid}/formatters/xml`** + - Returns `200 application/xml`, an `mdb:MD_Metadata` in the ISO 19115-3 + namespaces: about 32 KB for a release, 56 KB for a series. + - No authentication is needed for public records. + - This is the route Extract uses. +- **`GET {SDI_API_URL}/srv/api/site`** reports `EEA geospatial data catalogue`, + version `4.4.9`. +- **A series record lists its releases** as `mri:associatedResource` with + `associationType = isComposedOf` (21 of them for bathing water). + - A release record doesn't point back to its series. + - GeoNetwork's `/related?type=children` and a `parentUuid` search both + return nothing. +- **`POST {SDI_API_URL}/srv/api/search/records/_search`** with + `{"query": {"terms": {"uuid": [...]}}}` returns each record's: + - `cl_status` (`superseded` or absent) + - `resourceEdition` + - `publicationDateForResource` + - `resourceTemporalExtentDetails` + - `resourceTitleObject` + + Its `*DateRange` fields are shifted to UTC (31 December reads + `…-12-30T23:00Z`). Use the plain-date fields instead. + +## DDS interface assumptions (confirm with the DDS team) + +This repo only calls DDS's ingest API. Extract needs its document endpoints. +The plan assumes the following. `dds_documents/client.py` is written against +these assumptions, each held in a module constant, and is adjusted once the DDS +team confirms them: + +- **Upload:** `PUT {DDS_BASE_URL}/api/v1/files/{path}`, where the request body + is the file, stored as a document (never promoted to a Dremio table). + Requires write access on the parent folder. +- **Overwrite is explicit.** An existing file is rejected unless the request + says to overwrite it. Assumed: `?overwrite=true|false`, and **409** on + conflict (`DocumentExistsError`). +- **Download:** `GET {DDS_BASE_URL}/api/v1/files/{path}`, used to compare with + what's already there (404 when absent). Listing is assumed to be + `GET /api/v1/files?prefix={folder}`, returning paths. +- **A `metadata/` document folder at each catalog node.** What `{path}` is for + a dataset's metadata folder: the catalog path joined with `/`, then + `metadata/`, then the file name? Is the path case-folded? Does it include + the Dremio source prefix? Is the folder created on first upload? +- **Authentication** is the same `Authorization: Bearer ` the + ingest client sends. +- **Size limit** for an upload. ISO records are tens of KB, so any sane limit + is fine, but the client should report it clearly. + +## Design + +### Layout + +``` +src/eea_datalakehouse/ + sdi/ + __init__.py # public API: SdiController, SdiMetadata, PushResult, errors + config.py # SDI_API_URL (+ optional SDI_USERNAME / SDI_PASSWORD), from env + catalogue.py # read-only SDI REST client: get_xml(), search(), site() + iso.py # ISO 19115-3 read: uuid, title, edition, dates, hierarchy level + controller.py # SdiController.get_xml(), push_to_dds(), resolve_series() + errors.py # SdiError, SdiNotFound, SdiAuthError, UuidMismatch, NotIso19115_3, NotCurrentError + dds_documents/ + __init__.py + client.py # DocumentsClient: put(), get(), exists(), list() +``` + +- **`dds_documents` is its own package**, next to `dds_ingestion`, and the XML + reaches DDS through the **DDS REST API, the same way `IngestClient` does**: + + | | `IngestClient` today | `DocumentsClient` | + |---|---|---| + | Base URL | `DDS_BASE_URL` | same | + | Auth | `Authorization: Bearer ` (`_DREMIO_PWD`), never logged or in `repr` | same `DremioCreds` | + | Shape | one method per endpoint, typed results in `models.py` | same | + | Errors | `IngestApiError(status, message, where=…)`, error body capped | `DocumentsApiError`, same shape | + | Transport | one `httpx.Client`, context manager | same | + | Tests | `respx` | `respx` | + + `DremioCreds` and the `resolve_*` helpers move out of + `dds_ingestion/credentials.py` into a shared module, so both clients import + them from one place. + + **The XML is a document upload, not an ingest** (item 3 in + [Decide first](#decide-first)). + - Ingest (`/api/v1/ingest/*`) creates datasets in the catalog: `commit` + registers the upload as a Dremio table, and `DataFormat` is + `parquet | csv | json`. + - The SDI controller only moves metadata, so it never calls the ingest + API. `DocumentsClient` shares `IngestClient`'s conventions, not its + endpoints. +- **Transport is `httpx`**, like `dds_ingestion`, so everything is tested with + `respx` as the rest of the suite is. +- **ISO parsing uses `xml.etree.ElementTree`** with explicit namespaces. The + controller reads a handful of elements and doesn't need `lxml`. +- **Configuration comes from environment variables**, as `DDS_BASE_URL` does + today: + - `SDI_API_URL` (required). + - `SDI_USERNAME` / `SDI_PASSWORD` (optional), sent as basic auth, only for + reading non-public records. GeoNetwork 4.4 has no API-key mechanism. + - `DDS_BASE_URL`, `_DREMIO_USER`, `_DREMIO_PWD`, as today. + + All three are in `debugger/.env.example`. + +### Public API + +Two public methods, as built. The UUID always comes from the calling code: + +```python +from eea_datalakehouse.sdi import SdiController + +with SdiController.from_env() as sdi: # SDI_API_URL; DDS settings read on first push + metadata = sdi.get_xml("070d9baa-448d-4168-8514-7dadb3ad876d") + metadata.record # IsoRecord(uuid, title, edition, hierarchy_level, created, revised, children) + + res = sdi.push_to_dds( + metadata, + "catalog/water_management_resources/bathing_water/bwd", # the dataset's DDS path + # folder="metadata" (the default last segment) + # force=False (True replaces a DDS copy that was edited) + ) +res.dds_path # "catalog/.../bwd/metadata/070d9baa-448d-4168-8514-7dadb3ad876d.xml" +res.action # "uploaded" | "unchanged" | "replaced" +``` + +The full DDS path above the `metadata` folder is still to be agreed; +`push_to_dds` takes it `/`-joined from the caller, so nothing in the library +fixes it yet. + +**The destination is the `metadata` folder by default** (item 1 in +[Decide first](#decide-first)): + +- The file goes to `{dds_path}/metadata/{uuid}.xml` (`metadata_path()`). +- `folder=` overrides the sub-folder name. It is a keyword argument with the + default `"metadata"`, held in one constant (`DEFAULT_FOLDER`), so it is + changed in one place. +- The name is written lowercase. DDS may fold case (see + [DDS interface assumptions](#dds-interface-assumptions-confirm-with-the-dds-team)), + so `METADATA` and `metadata` must resolve to the same folder. + +### Extract (SDI → DDS) + +1. **Fetch** `GET {SDI_API_URL}/srv/api/records/{uuid}/formatters/xml`, for + the `uuid` the calling code passed in. + - A 404 raises `SdiNotFound`. + - A 401/403 on a non-public record raises `SdiAuthError`, with the hint to + set `SDI_USERNAME` / `SDI_PASSWORD`. +2. **Check that it's the right record.** `iso.py` parses the XML: + - The root must be `mdb:MD_Metadata`, otherwise `NotIso19115_3`. That covers + a 19139 `gmd:MD_Metadata` or an HTML error page. + - `mdb:metadataIdentifier/…/mcc:code` must equal the requested UUID, + otherwise `UuidMismatch`. +3. **Name the file** `{folder}/{uuid}.xml` under the dataset's DDS path, where + `folder` defaults to `metadata`. + - The UUID in both the file name and the content is what ties a DDS file to + its SDI record, so no separate link has to be stored. + - One file per record means a new release (new UUID) lands next to the old + one instead of overwriting it. + - Whether to keep only the current release is in + [Decide first](#decide-first). +4. **Compare** with what's already in DDS (a 404 means new). + - The bytes are equal: `unchanged`, no upload. + - They differ and the SDI record's revision date (`mdb:dateInfo`) is newer: + upload with overwrite and report `replaced`. + - They differ and the DDS copy is not older (newer, same date, undated or + not ISO at all), meaning someone edited it in DDS: raise + `DdsCopyConflict`. Extract never silently overwrites local edits. + `force=True` overrides. +5. **Upload** it as `application/xml` through `DocumentsClient.put()`. +6. **Return** a `PushResult(dds_path, action, uuid)`. Provenance on the wiki is + the caller's choice: `metadata.provenance_tags()` returns + `sdi_record_uuid` / `sdi_edition` / `sdi_date_stamp` in the `tags=` shape + of the existing `catalog.setmeta2wiki()`; the controller writes nothing to + Dremio itself. + +### Series resolution (Extract by series) + +A dataset is often a **series**, with one record per release. Letting Extract +take a series UUID means callers don't have to track which release is +current: + +- **`resolve_series(series_uuid)`**: + - Read the series XML and collect its `isComposedOf` children. + - Make one `_search` call for their status. + - The current release is the one child that isn't `superseded`. + - If there are zero or several, raise `NotCurrentError` listing each + candidate's UUID, title, edition, publication date and status, rather + than guessing by date. +- The caller then passes the returned UUID to `get_xml()`, so the UUID + that is extracted is still one the calling code handed over. +- **Live check (2026-10-01):** the bathing water series `c3858959…` has + **two** releases that are not superseded (`070d9baa…`, 2025 v1.0, and + `6b7b7f62…`, "Bathing water quality in Europe", 2020), so + `resolve_series` raises `NotCurrentError` for it. That is the intended + behaviour; pass the release UUID directly. + +### Notebook surface + +Two magics in `notebook/magics.py`, next to `%catalog` / `%ingest`, one session +per kernel each, built on first use from the kernel environment: + +- **`%sdi`** (`SdiSession`) reads SDI: `get_xml`, `resolve_series`. +- **`%metadata`** (`MetadataSession`) pushes metadata files to DDS: + `push_to_dds`, `dds_base_url`. + +``` +%sdi help +%sdi get_xml("070d9baa-448d-4168-8514-7dadb3ad876d") +%sdi resolve_series("c3858959-90da-4c1b-b9ca-492db0e514df") + +%metadata help +%metadata dds_base_url() +%metadata push_to_dds("catalog/water_management_resources/bathing_water/bwd") +``` + +- Each call runs immediately, like `%catalog`; there is nothing to commit. +- `%metadata push_to_dds` uploads the record `%sdi get_xml` fetched last + (looked up at push time) unless `metadata=` is given. +- Every failure prints one `sdi error: …` / `metadata error: …` line instead + of a traceback. +- `DDS_BASE_URL` is read, when `%metadata` builds its session, from the first + `.env` in the notebook's folder or a parent; the kernel environment's value + applies only if no `.env` sets it. `%metadata dds_base_url()` shows which is + in use. +- The DDS upload's Dremio identity is `_DREMIO_USER` / `_DREMIO_PWD` (as + `%ingest`), falling back to `DREMIO_USERNAME` / `DREMIO_TOKEN` (as + `%catalog`). + +## Steps + +Each step lands with its tests, mypy clean and ruff clean. + +0. **Confirm the [DDS interface](#dds-interface-assumptions-confirm-with-the-dds-team)** + with the DDS team. +1. ✅ **`sdi/config.py`, `sdi/catalogue.py`, `sdi/iso.py`.** + - Test records are built in `tests/sdi/conftest.py` with the same structure + the live catalogue serves. + - `respx` mocks the HTTP calls. + - One live test, skipped unless `SDI_NETWORK_TESTS=1`. +2. ✅ **`dds_documents/client.py`**, built like `IngestClient`: `put` / `get` / + `exists` / `list` over the document endpoints. It makes no ingest calls. + - Written against the assumed contract; each endpoint, the `overwrite` + query parameter and the 409 conflict are module constants, to adjust + once step 0 confirms them. The size limit and path case wait for step 0. + - The credential code stays in `dds_ingestion/credentials.py` for now and + is imported from there. +3. ✅ **`SdiController.get_xml()` / `push_to_dds()`**, all of Extract, plus + `resolve_series()`. + - Cases: new, unchanged, replaced, refused because the DDS copy was edited, + `force`, `UuidMismatch`, not ISO 19115-3, and zero or several current + releases. + - `debugger/debug_run.py` has `run_sdi_extract()` for the bathing water + release (extracts for real; uploads only with `DRY_RUN = False`). +4. **Prove Extract against a DDS test instance.** Extract the bathing water + release, and check the file appears in the dataset's `metadata/` folder. +5. ✅ **`%sdi` and `%metadata` magics** over `sdi/session.py`'s `SdiSession` + (`get_xml`, `resolve_series`) and `MetadataSession` (`push_to_dds`, which + defaults to the last `%sdi` record, and `dds_base_url`). Built on first use + from the kernel environment like `%catalog` / `%ingest`. + Worked example: `debugger/sdi_session_example.ipynb`. The README + section, Layout rows, docstrings and `.env.example` entries are done. +6. **Dependencies.** + - Drop `lxml`, since `iso.py` uses ElementTree. + - Drop `beautifulsoup4`, which nothing here needs. + - Correct `xlrd`'s comment. + +## Decide first + +1. ✅ **Decided (2026-09-30): the destination is the `metadata` folder by + default**, i.e. `{dds_path}/metadata/{uuid}.xml`. + - The calling code passes `dds_path`, so it picks the catalog node + (dataflow, table, and so on). The full path convention is still to be + agreed. + - `folder=` overrides the sub-folder. + + Still open, and it doesn't block building: one file per SDI record + (`{uuid}.xml`, several releases side by side, the default in this plan), or + a single `current.xml` per dataset? +2. ✅ **Decided (2026-09-30): the SDI UUID is always a parameter from the + calling code.** + - `get_xml(uuid)` / `resolve_series(series_uuid)` take it explicitly. + - The library never looks it up: not from the catalog, not from a wiki + key, not from a file name. + - `provenance_tags()` only records the UUID for people to read. Nothing + reads it back to pick a record. +3. ✅ **Decided (2026-09-30): the XML goes to DDS as a document upload, not + through ingest.** + - Ingest is for creating datasets in the catalog. The SDI controller is + only about metadata. + - `DocumentsClient` uses DDS's document endpoints (`PUT` / + `GET /api/v1/files/{path}`). Their exact contract is still to confirm, + see + [DDS interface assumptions](#dds-interface-assumptions-confirm-with-the-dds-team). + +## Not in this plan + +- **Publication to SDI (DDS → SDI).** Deferred (2026-10-01). When it is taken + up, it needs its own decisions first: SDI write credentials (an editor + account; an unauthenticated `PUT /srv/api/records` returns 403, and + GeoNetwork 4 also needs its XSRF cookie/header pair), EEA SDI's editorial + workflow, and who may publish. The Extract design leaves room for it: + `catalogue.py` can gain a write side, and the `{uuid}.xml` naming already + ties each DDS file to its SDI record. +- **Editing metadata.** Editing the XML itself is done in DDS by people. The + controller only moves validated ISO 19115-3 into DDS. +- **Full schema validation** against the ISO 19115-3 XSDs or INSPIRE rules. + GeoNetwork has already validated what it serves. +- **ISO 19139 (`gmd:`) records.** The controller refuses them rather than + converting them. +- **Downloading the data** behind an SDI record. This plan covers metadata + only. +- **Search and harvesting** by keyword or title. +- **CSW and OGC API Records.** The GeoNetwork REST API is enough. +- **Changes to DDS.** The plan uses DDS's existing document endpoints, as + confirmed. + +## Sources + +- This repo: + - `src/eea_datalakehouse/dds_ingestion/client.py` and `credentials.py` + - `src/eea_datalakehouse/catalog/operations.py` (`setmeta2wiki`, + `getmetafromwiki`) + - `src/eea_datalakehouse/notebook/magics.py` + - `docs/dev-notes.md` + - `pyproject.toml` +- The live SDI catalogue, 2026-09-30 (see + [Verified](#verified-against-the-live-sdi-catalogue-2026-09-30)) +- DDS: the assumptions above, to be confirmed with the DDS team diff --git a/src/eea_datalakehouse/dds_documents/__init__.py b/src/eea_datalakehouse/dds_documents/__init__.py new file mode 100644 index 0000000..c88fac6 --- /dev/null +++ b/src/eea_datalakehouse/dds_documents/__init__.py @@ -0,0 +1,16 @@ +"""eea_datalakehouse.dds_documents — store plain files (documents) in DDS folders. + +Unlike :mod:`eea_datalakehouse.dds_ingestion`, nothing here creates a dataset: +a document is uploaded as-is and never becomes a Dremio table:: + + from eea_datalakehouse.dds_documents import DocumentsClient + + with DocumentsClient.from_env() as docs: + docs.put("a/b/metadata/x.xml", xml_bytes, content_type="application/xml") +""" + +from __future__ import annotations + +from .client import DocumentExistsError, DocumentsApiError, DocumentsClient + +__all__ = ["DocumentExistsError", "DocumentsApiError", "DocumentsClient"] diff --git a/src/eea_datalakehouse/dds_documents/client.py b/src/eea_datalakehouse/dds_documents/client.py new file mode 100644 index 0000000..1b6e0b4 --- /dev/null +++ b/src/eea_datalakehouse/dds_documents/client.py @@ -0,0 +1,190 @@ +"""Thin HTTP client for the DDS document endpoints. + +Documents are plain files stored in a DDS folder (an ISO 19115-3 XML in a +dataset's ``metadata`` folder, for instance). They are never registered as +Dremio tables, so this client makes **no ingest calls**: it shares +:class:`~eea_datalakehouse.dds_ingestion.client.IngestClient`'s conventions +(base URL, Bearer PAT, one method per endpoint, errors naming the call), not +its endpoints. + +=============================== ============================================= +:meth:`DocumentsClient.put` ``PUT /api/v1/files/{path}`` +:meth:`~DocumentsClient.get` ``GET /api/v1/files/{path}`` +:meth:`~DocumentsClient.exists` ``GET /api/v1/files/{path}`` +:meth:`~DocumentsClient.list` ``GET /api/v1/files?prefix={folder}`` +=============================== ============================================= + +The endpoint shapes below are **interface assumptions** still to confirm with +the DDS team (see ``docs/sdi-integration-plan.md``). Each one is a module +constant so a confirmed contract is a one-line change. +""" + +from __future__ import annotations + +from types import TracebackType +from typing import Any +from urllib.parse import quote + +import httpx + +from ..dds_ingestion.credentials import DremioCreds, load_base_url, load_creds + +DEFAULT_TIMEOUT = 60.0 +# Assumed DDS document contract — to confirm with the DDS team. +FILES_ENDPOINT = "/api/v1/files/{path}" +LIST_ENDPOINT = "/api/v1/files" +OVERWRITE_PARAM = "overwrite" +# Error bodies are quoted back to the caller; cap them so an HTML error page +# does not bury the status line it came with. +_MAX_ERROR_CHARS = 500 + + +class DocumentsApiError(RuntimeError): + """A DDS document call returned a non-success status. + + ``where`` names the call that failed (e.g. ``PUT /api/v1/files/a/b.xml``). + """ + + def __init__(self, status_code: int, message: str, *, where: str | None = None) -> None: + super().__init__( + f"DDS documents API error {status_code}{f' on {where}' if where else ''}: {message}" + ) + self.status_code = status_code + self.message = message + self.where = where + + +class DocumentExistsError(DocumentsApiError): + """The file already exists and the upload did not ask to overwrite it (409).""" + + +class DocumentsClient: + """Client for the DDS document endpoints.""" + + def __init__( + self, + base_url: str, + creds: DremioCreds, + *, + timeout: float = DEFAULT_TIMEOUT, + http_client: httpx.Client | None = None, + ) -> None: + self._base_url = base_url.rstrip("/") + self._owns_client = http_client is None + self._http = http_client or httpx.Client(timeout=timeout) + self._auth_headers = {"Authorization": f"Bearer {creds.password}"} + + @classmethod + def from_env( + cls, env: dict[str, str] | None = None, *, timeout: float = DEFAULT_TIMEOUT + ) -> DocumentsClient: + """Build from ``DDS_BASE_URL`` and ``_DREMIO_USER`` / ``_DREMIO_PWD``.""" + return cls(load_base_url(env), load_creds(env), timeout=timeout) + + def __repr__(self) -> str: + # Deliberately omits _auth_headers so the PAT never reaches output. + return f"DocumentsClient(base_url={self._base_url!r})" + + # -- lifecycle -------------------------------------------------------- + + def close(self) -> None: + if self._owns_client: + self._http.close() + + def __enter__(self) -> DocumentsClient: + return self + + def __exit__( + self, + exc_type: type[BaseException] | None, + exc: BaseException | None, + tb: TracebackType | None, + ) -> None: + self.close() + + # -- API endpoints ---------------------------------------------------- + + def put( + self, + path: str, + data: bytes, + *, + content_type: str = "application/octet-stream", + overwrite: bool = False, + ) -> None: + """Store ``data`` as the document at ``path``. + + Without ``overwrite`` an existing document is left alone and + :class:`DocumentExistsError` is raised. + """ + endpoint = self._file_endpoint(path) + resp = self._http.put( + f"{self._base_url}{endpoint}", + content=data, + params={OVERWRITE_PARAM: "true" if overwrite else "false"}, + headers={**self._auth_headers, "Content-Type": content_type}, + ) + if resp.status_code == 409: + raise DocumentExistsError(409, _error_message(resp), where=f"PUT {endpoint}") + _raise_for_status(resp, f"PUT {endpoint}") + + def get(self, path: str) -> bytes | None: + """The document's bytes, or ``None`` if there is none at ``path``.""" + endpoint = self._file_endpoint(path) + resp = self._http.get(f"{self._base_url}{endpoint}", headers=self._auth_headers) + if resp.status_code == 404: + return None + _raise_for_status(resp, f"GET {endpoint}") + return resp.content + + def exists(self, path: str) -> bool: + """Whether a document is stored at ``path``.""" + return self.get(path) is not None + + def list(self, folder: str) -> list[str]: + """Paths of the documents under ``folder``.""" + prefix = _clean_path(folder) + resp = self._http.get( + f"{self._base_url}{LIST_ENDPOINT}", + params={"prefix": prefix}, + headers=self._auth_headers, + ) + _raise_for_status(resp, f"GET {LIST_ENDPOINT}?prefix={prefix}") + payload: Any = resp.json() + rows = payload.get("files", []) if isinstance(payload, dict) else payload + return [row["path"] if isinstance(row, dict) else str(row) for row in rows] + + # -- internals -------------------------------------------------------- + + @staticmethod + def _file_endpoint(path: str) -> str: + return FILES_ENDPOINT.format(path=quote(_clean_path(path), safe="/")) + + +def _clean_path(path: str) -> str: + cleaned = path.strip().strip("/") + if not cleaned: + raise ValueError("DDS document path must not be empty") + if any(part in ("", ".", "..") for part in cleaned.split("/")): + raise ValueError(f"invalid DDS document path {path!r}") + return cleaned + + +def _error_message(resp: httpx.Response) -> str: + """The ``message`` of a DDS ``{error, message}`` body, else the capped text.""" + message = resp.text + try: + payload = resp.json() + if isinstance(payload, dict): + message = str(payload.get("message") or payload.get("error") or message) + except ValueError: + pass + message = message.strip() + if len(message) > _MAX_ERROR_CHARS: + message = message[:_MAX_ERROR_CHARS] + "… (truncated)" + return message or f"(empty {resp.status_code} response body)" + + +def _raise_for_status(resp: httpx.Response, where: str) -> None: + if not resp.is_success: + raise DocumentsApiError(resp.status_code, _error_message(resp), where=where) diff --git a/src/eea_datalakehouse/notebook/__init__.py b/src/eea_datalakehouse/notebook/__init__.py index 13578ac..6de4e0b 100644 --- a/src/eea_datalakehouse/notebook/__init__.py +++ b/src/eea_datalakehouse/notebook/__init__.py @@ -1,10 +1,10 @@ -"""The notebook-facing surface — `%catalog`/`%ingest` magics over -`CatalogSession`/`IngestSession` (see `docs/notebook-facade-for-data- -scientists.md`). Requires IPython, which the rest of this package does not -— install the `notebook` extra (`pip install "EEADataLakehouse[notebook]"`) -to get it. +"""The notebook-facing surface — `%catalog`/`%ingest`/`%sdi`/`%metadata` magics +over `CatalogSession`/`IngestSession`/`SdiSession`/`MetadataSession` (see +`docs/notebook-facade-for-data-scientists.md`). Requires IPython, which the +rest of this package does not — install the `notebook` extra +(`pip install "EEADataLakehouse[notebook]"`) to get it. -Importing this package inside a running IPython shell registers both magics +Importing this package inside a running IPython shell registers every magic as a side effect — a custodian only ever needs:: import eea_datalakehouse.notebook diff --git a/src/eea_datalakehouse/notebook/magics.py b/src/eea_datalakehouse/notebook/magics.py index 04c78aa..54a7059 100644 --- a/src/eea_datalakehouse/notebook/magics.py +++ b/src/eea_datalakehouse/notebook/magics.py @@ -1,10 +1,10 @@ -"""`%catalog`/`%%catalog` and `%ingest` — thin syntactic sugar over -`CatalogSession`/`IngestSession` (see +"""`%catalog`/`%%catalog`, `%ingest`, `%sdi` and `%metadata` — thin syntactic +sugar over `CatalogSession`/`IngestSession`/`SdiSession`/`MetadataSession` (see `docs/notebook-facade-for-data-scientists.md`, "Two magics, two sessions"). Available the moment this package is imported inside IPython — `import -eea_datalakehouse.notebook` (see that package's `__init__.py`) registers both -magics as a side effect, so nothing needs loading explicitly. `%load_ext +eea_datalakehouse.notebook` (see that package's `__init__.py`) registers every +magic as a side effect, so nothing needs loading explicitly. `%load_ext eea_datalakehouse.notebook.magics` still works too (and is the only option outside IPython's auto-import path, e.g. a config that imports this module directly) — `load_ipython_extension` below is a no-op if the auto-import @@ -20,6 +20,16 @@ data_format="parquet", table_name="water_temperature") %ingest commit(retry=True) + %sdi get_xml("070d9baa-448d-4168-8514-7dadb3ad876d") + %metadata push_to_dds("catalog/water_management_resources/bathing_water/bwd") + +`%sdi` (`eea_datalakehouse.sdi.session.SdiSession`) only reads the SDI +catalogue; `%metadata` (`MetadataSession`, same module) pushes metadata files +to DDS. Both run each call immediately, like `%catalog`, and are built on first +use, like `%ingest`'s session, from the kernel environment (see that module's +docstring for the variables). `%metadata push_to_dds` uploads the record +`%sdi get_xml` fetched last, so it needs only the DDS path. + `%catalog` executes each call immediately — there's no queue and no separate commit step to remember. A `CatalogSession` still sits underneath it (to keep the "current path" context, and the same retry-on-`EngineStartingError`/ @@ -87,6 +97,7 @@ from ..catalog import Catalog from ..catalog.session import CatalogSession, CatalogSessionError from ..dds_ingestion.session import IngestSession, IngestSessionError +from ..sdi.session import MetadataSession, MetadataSessionError, SdiSession, SdiSessionError _USAGE = { "catalog": ( @@ -97,6 +108,11 @@ '%ingest ingest(folder="./data", target_catalog_path="a.b", data_format="parquet")' " | %ingest commit" ), + "sdi": '%sdi get_xml("") | %sdi resolve_series("")', + "metadata": ( + '%metadata push_to_dds("a/b/c") (uploads the last %sdi get_xml result)' + " | %metadata dds_base_url()" + ), } # Shared between set_wiki's _CATALOG_HELP entry below and set_tags'/ @@ -233,6 +249,42 @@ ("ingest", "Queue one folder ingest (see FolderIngest for what each argument means)."), ("commit", "Run every queued ingest in order; stops at the first failure."), ] +_SDI_HELP = [ + ( + "get_xml", + "Download the ISO 19115-3 XML of SDI record uuid, check it really is that record, " + "and remember it for %metadata push_to_dds. Runs immediately; needs no DDS settings.", + ), + ( + "resolve_series", + "Return the UUID of the one release of series_uuid that is not superseded; raises, " + "listing every release, if there are zero or several.", + ), +] +_SDI_HELP_NOTE = ( + "Environment: SDI_API_URL (empty means the public EEA catalogue; optional " + "SDI_USERNAME/SDI_PASSWORD for non-public records). Pushing to DDS is %metadata's job." +) +_METADATA_HELP = [ + ( + "push_to_dds", + "Upload metadata (default: the last %sdi get_xml result) to " + "dds_path/folder/.xml in DDS — a document upload, never an ingest. Same bytes " + "already there: unchanged; an older copy: replaced; a copy edited in DDS: refused " + "unless force=True.", + ), + ( + "dds_base_url", + "Show the DDS URL push_to_dds uses and where it came from (a .env file or the kernel " + "environment).", + ), +] +_METADATA_HELP_NOTE = ( + "Environment: DDS_BASE_URL — read from the first .env found in the notebook's folder " + "or a parent when %metadata starts, else from the kernel environment — and a Dremio " + "identity: _DREMIO_USER/_DREMIO_PWD (as %ingest) or DREMIO_USERNAME/DREMIO_TOKEN " + "(as %catalog)." +) def _build_catalog_session() -> CatalogSession: @@ -246,6 +298,18 @@ def _build_catalog_session() -> CatalogSession: return CatalogSession(Catalog(base_url, token, username=username)) +def _build_sdi_session() -> SdiSession: + # Nothing to check up front: SDI_API_URL defaults to the public catalogue. + return SdiSession() + + +def _build_metadata_session(last: Any) -> MetadataSession: + # MetadataSession reads DDS_BASE_URL from the nearest .env right here, when + # it is built; a missing DDS setting or Dremio identity prints as a + # friendly "metadata error" on the first push_to_dds. + return MetadataSession(last=last) + + def _is_help(line: str) -> bool: return line.strip() in ("help", "help()") @@ -412,12 +476,14 @@ def _on_msg(msg: dict[str, Any]) -> None: @magics_class class EEALakehouseMagics(Magics): - """Registers `%catalog` and `%ingest` — see this module's docstring.""" + """Registers `%catalog`, `%ingest`, `%sdi` and `%metadata` — see this module's docstring.""" def __init__(self, shell: Any) -> None: super().__init__(shell) self._catalog_session: CatalogSession | None = None self._ingest_session: IngestSession | None = None + self._sdi_session: SdiSession | None = None + self._metadata_session: MetadataSession | None = None self._pending_context: str | None = None _register_context_comm(shell, self) @@ -523,6 +589,46 @@ def ingest(self, line: str) -> Any: ) return None if result is _DISPATCH_FAILED else result + @line_magic + def sdi(self, line: str) -> Any: + """`%sdi` — read ISO 19115-3 metadata from SDI, each call runs immediately:: + + %sdi get_xml("070d9baa-448d-4168-8514-7dadb3ad876d") + %sdi resolve_series("c3858959-90da-4c1b-b9ca-492db0e514df") + """ + if _is_help(line): + _print_help_table(SdiSession, _SDI_HELP, "sdi", note=_SDI_HELP_NOTE) + return None + if self._sdi_session is None: + self._sdi_session = _build_sdi_session() + result = _dispatch(self._sdi_session, line, self.shell.user_ns, "sdi", SdiSessionError) + return None if result is _DISPATCH_FAILED else result + + @line_magic + def metadata(self, line: str) -> Any: + """`%metadata` — push metadata files to DDS, each call runs immediately:: + + %metadata push_to_dds("catalog/water_management_resources/bathing_water/bwd") + %metadata dds_base_url() + + `push_to_dds` uploads the record `%sdi get_xml` fetched last unless + `metadata=` is given. + """ + if _is_help(line): + _print_help_table( + MetadataSession, _METADATA_HELP, "metadata", note=_METADATA_HELP_NOTE + ) + return None + if self._metadata_session is None: + # Looked up at push time, so a get_xml run after %metadata started still counts. + self._metadata_session = _build_metadata_session( + lambda: self._sdi_session.last if self._sdi_session is not None else None + ) + result = _dispatch( + self._metadata_session, line, self.shell.user_ns, "metadata", MetadataSessionError + ) + return None if result is _DISPATCH_FAILED else result + def load_ipython_extension(ipython: Any) -> None: # A no-op if `eea_datalakehouse.notebook`'s own auto-import already diff --git a/src/eea_datalakehouse/sdi/__init__.py b/src/eea_datalakehouse/sdi/__init__.py new file mode 100644 index 0000000..63baaf2 --- /dev/null +++ b/src/eea_datalakehouse/sdi/__init__.py @@ -0,0 +1,63 @@ +"""eea_datalakehouse.sdi — ISO 19115-3 metadata from the EEA SDI catalogue to DDS. + +Typical use (``SDI_API_URL``, ``DDS_BASE_URL`` and the Dremio credentials come +from the kernel env):: + + from eea_datalakehouse.sdi import SdiController + + with SdiController.from_env() as sdi: + metadata = sdi.get_xml("070d9baa-448d-4168-8514-7dadb3ad876d") + result = sdi.push_to_dds(metadata, "water_management_resources/bathing_water/bwd") + result.dds_path # ".../bwd/metadata/070d9baa-448d-4168-8514-7dadb3ad876d.xml" + result.action # "uploaded" | "unchanged" | "replaced" +""" + +from __future__ import annotations + +from .catalogue import SdiCatalogue +from .config import DEFAULT_SDI_API_URL, SdiConfig +from .controller import ( + DEFAULT_FOLDER, + PushResult, + SdiController, + SdiMetadata, + metadata_path, +) +from .errors import ( + DdsCopyConflict, + NotCurrentError, + NotIso19115_3, + SdiApiError, + SdiAuthError, + SdiError, + SdiNotFound, + SeriesCandidate, + UuidMismatch, +) +from .iso import IsoRecord +from .session import MetadataSession, MetadataSessionError, SdiSession, SdiSessionError + +__all__ = [ + "DEFAULT_FOLDER", + "DEFAULT_SDI_API_URL", + "DdsCopyConflict", + "IsoRecord", + "MetadataSession", + "MetadataSessionError", + "NotCurrentError", + "NotIso19115_3", + "PushResult", + "SdiApiError", + "SdiAuthError", + "SdiCatalogue", + "SdiConfig", + "SdiController", + "SdiError", + "SdiMetadata", + "SdiNotFound", + "SdiSession", + "SdiSessionError", + "SeriesCandidate", + "UuidMismatch", + "metadata_path", +] diff --git a/src/eea_datalakehouse/sdi/catalogue.py b/src/eea_datalakehouse/sdi/catalogue.py new file mode 100644 index 0000000..5e93011 --- /dev/null +++ b/src/eea_datalakehouse/sdi/catalogue.py @@ -0,0 +1,155 @@ +"""Read-only client for the SDI (GeoNetwork 4) REST API. + +====================================== ===================================================== +:meth:`SdiCatalogue.get_xml` ``GET {SDI_API_URL}/srv/api/records/{uuid}/formatters/xml`` +:meth:`~SdiCatalogue.search_by_uuid` ``POST {SDI_API_URL}/srv/api/search/records/_search`` +:meth:`~SdiCatalogue.site` ``GET {SDI_API_URL}/srv/api/site`` +====================================== ===================================================== +""" + +from __future__ import annotations + +from types import TracebackType +from typing import Any +from urllib.parse import quote + +import httpx + +from .config import SdiConfig +from .errors import SdiApiError, SdiAuthError, SdiNotFound, SeriesCandidate + +DEFAULT_TIMEOUT = 60.0 +_MAX_ERROR_CHARS = 500 +# Fields asked of the search index. Its *DateRange fields are shifted to UTC +# (31 December reads …-12-30T23:00Z), so only plain-date fields are used. +_SEARCH_FIELDS = [ + "uuid", + "resourceTitleObject", + "resourceEdition", + "publicationDateForResource", + "cl_status", +] + + +class SdiCatalogue: + """Read-only access to one GeoNetwork catalogue.""" + + def __init__( + self, + config: SdiConfig | None = None, + *, + timeout: float = DEFAULT_TIMEOUT, + http_client: httpx.Client | None = None, + ) -> None: + self._config = config or SdiConfig() + self._api = f"{self._config.api_url.rstrip('/')}/srv/api" + self._owns_client = http_client is None + self._http = http_client or httpx.Client(timeout=timeout, follow_redirects=True) + + @classmethod + def from_env( + cls, env: dict[str, str] | None = None, *, timeout: float = DEFAULT_TIMEOUT + ) -> SdiCatalogue: + return cls(SdiConfig.from_env(env), timeout=timeout) + + def __repr__(self) -> str: + return f"SdiCatalogue(api_url={self._config.api_url!r})" + + # -- lifecycle -------------------------------------------------------- + + def close(self) -> None: + if self._owns_client: + self._http.close() + + def __enter__(self) -> SdiCatalogue: + return self + + def __exit__( + self, + exc_type: type[BaseException] | None, + exc: BaseException | None, + tb: TracebackType | None, + ) -> None: + self.close() + + # -- API endpoints ---------------------------------------------------- + + def get_xml(self, uuid: str) -> bytes: + """The record's ISO XML, exactly as SDI serves it.""" + path = f"/records/{quote(uuid, safe='')}/formatters/xml" + resp = self._http.get( + f"{self._api}{path}", + headers={"Accept": "application/xml"}, + **self._auth(), + ) + self._raise_for_status(resp, f"GET /srv/api{path}", uuid=uuid) + return resp.content + + def search_by_uuid(self, uuids: list[str]) -> list[SeriesCandidate]: + """Title, edition, publication date and status of each record in ``uuids``.""" + if not uuids: + return [] + body = { + "query": {"terms": {"uuid": uuids}}, + "size": len(uuids), + "_source": _SEARCH_FIELDS, + } + resp = self._http.post( + f"{self._api}/search/records/_search", + json=body, + headers={"Accept": "application/json"}, + **self._auth(), + ) + self._raise_for_status(resp, "POST /srv/api/search/records/_search") + hits = resp.json().get("hits", {}).get("hits", []) + return [_candidate(hit.get("_source", {})) for hit in hits] + + def site(self) -> dict[str, Any]: + """The catalogue's site description (name, version).""" + resp = self._http.get( + f"{self._api}/site", headers={"Accept": "application/json"}, **self._auth() + ) + self._raise_for_status(resp, "GET /srv/api/site") + result: dict[str, Any] = resp.json() + return result + + # -- internals -------------------------------------------------------- + + def _auth(self) -> dict[str, Any]: + """``auth=`` for a request: basic auth when configured, else none.""" + auth = self._config.auth + return {"auth": auth} if auth is not None else {} + + @staticmethod + def _raise_for_status(resp: httpx.Response, where: str, *, uuid: str | None = None) -> None: + if resp.is_success: + return + if resp.status_code == 404 and uuid is not None: + raise SdiNotFound(404, f"no SDI record with UUID {uuid!r}", where=where) + message = resp.text.strip() + if len(message) > _MAX_ERROR_CHARS: + message = message[:_MAX_ERROR_CHARS] + "… (truncated)" + message = message or f"(empty {resp.status_code} response body)" + if resp.status_code in (401, 403): + raise SdiAuthError( + resp.status_code, + f"{message} (record not public? set SDI_USERNAME / SDI_PASSWORD)", + where=where, + ) + raise SdiApiError(resp.status_code, message, where=where) + + +def _candidate(source: dict[str, Any]) -> SeriesCandidate: + title = source.get("resourceTitleObject") + published = source.get("publicationDateForResource") + statuses = source.get("cl_status") or [] + if isinstance(statuses, dict): + statuses = [statuses] + return SeriesCandidate( + uuid=str(source.get("uuid", "")), + title=title.get("default") if isinstance(title, dict) else title, + edition=source.get("resourceEdition"), + publication_date=str(published[0]) if isinstance(published, list) and published + else (str(published) if published else None), + status=next((s.get("key") for s in statuses if isinstance(s, dict)), None), + ) diff --git a/src/eea_datalakehouse/sdi/config.py b/src/eea_datalakehouse/sdi/config.py new file mode 100644 index 0000000..d67e8f1 --- /dev/null +++ b/src/eea_datalakehouse/sdi/config.py @@ -0,0 +1,50 @@ +"""SDI connection settings, read from the environment. + +``SDI_API_URL`` is the catalogue's base URL (the client appends +``/srv/api/...``); empty or unset means the public EEA catalogue. +``SDI_USERNAME`` / ``SDI_PASSWORD`` are optional and only needed to read +records that are not public. The password is never rendered by ``repr``. +""" + +from __future__ import annotations + +import os +from dataclasses import dataclass + +import httpx + +ENV_API_URL = "SDI_API_URL" +ENV_USERNAME = "SDI_USERNAME" +ENV_PASSWORD = "SDI_PASSWORD" # noqa: S105 — env var name, not a secret value +DEFAULT_SDI_API_URL = "https://sdi.eea.europa.eu/catalogue" + + +@dataclass(frozen=True) +class SdiConfig: + api_url: str = DEFAULT_SDI_API_URL + username: str | None = None + _password: str | None = None + + @classmethod + def from_env(cls, env: dict[str, str] | None = None) -> SdiConfig: + source = os.environ if env is None else env + url = (source.get(ENV_API_URL) or DEFAULT_SDI_API_URL).rstrip("/") + return cls( + api_url=url, + username=source.get(ENV_USERNAME) or None, + _password=source.get(ENV_PASSWORD) or None, + ) + + @property + def auth(self) -> httpx.BasicAuth | None: + if self.username and self._password: + return httpx.BasicAuth(self.username, self._password) + return None + + def __repr__(self) -> str: + return ( + f"SdiConfig(api_url={self.api_url!r}, username={self.username!r}, " + f"password={'***' if self._password else None})" + ) + + __str__ = __repr__ diff --git a/src/eea_datalakehouse/sdi/controller.py b/src/eea_datalakehouse/sdi/controller.py new file mode 100644 index 0000000..c66a85f --- /dev/null +++ b/src/eea_datalakehouse/sdi/controller.py @@ -0,0 +1,249 @@ +"""SDI metadata controller: SDI catalogue → DDS ``metadata`` folder. + +Two public methods do the work, and the calling code always passes the UUID: + +================================= ================================================ +:meth:`SdiController.get_xml` download a record's ISO 19115-3 XML from SDI +:meth:`SdiController.push_to_dds` upload it to ``{dds_path}/metadata/{uuid}.xml`` +================================= ================================================ + +plus :meth:`~SdiController.resolve_series`, which turns a series UUID into its +current release's UUID for callers that only know the series. + +The XML goes to DDS as a **document upload**, never through the ingest API: +ingest creates datasets in the catalog, this controller only moves metadata. +""" + +from __future__ import annotations + +from dataclasses import dataclass, field +from types import TracebackType +from typing import Literal + +from ..dds_documents import DocumentsClient +from . import iso +from .catalogue import SdiCatalogue +from .errors import DdsCopyConflict, NotCurrentError, NotIso19115_3 +from .iso import IsoRecord + +DEFAULT_FOLDER = "metadata" +XML_CONTENT_TYPE = "application/xml" + +PushAction = Literal["uploaded", "unchanged", "replaced"] + + +@dataclass(frozen=True) +class SdiMetadata: + """One ISO 19115-3 record as SDI served it. + + ``xml`` holds the exact bytes, so what lands in DDS is byte-for-byte what + SDI returned. ``record`` is the parsed summary (title, edition, dates). + """ + + uuid: str + xml: bytes = field(repr=False) + record: IsoRecord = field(repr=False) + + def __repr__(self) -> str: + return ( + f"SdiMetadata(uuid={self.uuid!r}, title={self.record.title!r}, " + f"edition={self.record.edition!r}, xml=<{len(self.xml)} bytes>)" + ) + + def provenance_tags(self) -> list[dict[str, str]]: + """Tags in :meth:`Catalog.setmeta2wiki`'s shape, recording which SDI record this is. + + For people reading the wiki only; nothing reads them back to pick a + record. + """ + stamp = self.record.date_stamp + values = { + "sdi_record_uuid": ("SDI record UUID", self.uuid), + "sdi_edition": ("SDI edition", self.record.edition or ""), + "sdi_date_stamp": ("SDI metadata date", stamp.isoformat() if stamp else ""), + } + return [ + {"tag_name": name, "tag_value": value, "tag_title": title} + for name, (title, value) in values.items() + ] + + +@dataclass(frozen=True) +class PushResult: + """What :meth:`SdiController.push_to_dds` did.""" + + dds_path: str + action: PushAction + uuid: str + + +class SdiController: + """Moves ISO 19115-3 metadata records from SDI into DDS. + + The DDS client is optional at construction: :meth:`get_xml` needs only + SDI, and :meth:`push_to_dds` builds a :class:`DocumentsClient` from the + environment (``DDS_BASE_URL``, ``_DREMIO_USER`` / ``_DREMIO_PWD``) on first + use if none was given. + """ + + def __init__( + self, + catalogue: SdiCatalogue | None = None, + documents: DocumentsClient | None = None, + ) -> None: + self._catalogue = catalogue or SdiCatalogue() + self._documents = documents + self._owns_documents = documents is None + + @classmethod + def from_env(cls, env: dict[str, str] | None = None) -> SdiController: + """SDI settings from ``SDI_API_URL`` (+ optional ``SDI_USERNAME`` / ``SDI_PASSWORD``). + + DDS settings are read later, by :meth:`push_to_dds`, so extracting + works without them. + """ + return cls(SdiCatalogue.from_env(env)) + + def __repr__(self) -> str: + return f"SdiController(catalogue={self._catalogue!r}, documents={self._documents!r})" + + @property + def catalogue(self) -> SdiCatalogue: + """The SDI client this controller reads from.""" + return self._catalogue + + # -- lifecycle -------------------------------------------------------- + + def close(self) -> None: + self._catalogue.close() + if self._documents is not None and self._owns_documents: + self._documents.close() + + def __enter__(self) -> SdiController: + return self + + def __exit__( + self, + exc_type: type[BaseException] | None, + exc: BaseException | None, + tb: TracebackType | None, + ) -> None: + self.close() + + # -- public API ------------------------------------------------------- + + def get_xml(self, uuid: str) -> SdiMetadata: + """Download the ISO 19115-3 XML of the SDI record ``uuid``. + + Raises :class:`SdiNotFound` if SDI has no such record, + :class:`SdiAuthError` if it is not public, :class:`SdiApiError` on any + other failure, :class:`NotIso19115_3` if the response is not an + ``mdb:MD_Metadata`` document and :class:`UuidMismatch` if its + identifier is not ``uuid``. + """ + uuid = _require_uuid(uuid) + xml = self._catalogue.get_xml(uuid) + try: + record = iso.check(xml, uuid) + except NotIso19115_3 as exc: + raise NotIso19115_3(f"SDI record {uuid!r}: {exc}") from exc + return SdiMetadata(uuid=uuid, xml=xml, record=record) + + def push_to_dds( + self, + metadata: SdiMetadata, + dds_path: str, + *, + folder: str = DEFAULT_FOLDER, + force: bool = False, + ) -> PushResult: + """Upload ``metadata`` to DDS as ``{dds_path}/{folder}/{uuid}.xml``. + + ``dds_path`` is the dataset's DDS location, ``/``-joined; ``folder`` is + the last path segment, ``metadata`` by default, always lowercase. + + Compared with what DDS already holds there: + + - nothing: upload, ``"uploaded"``; + - the same bytes: no upload, ``"unchanged"``; + - different bytes and an older metadata date: overwrite, ``"replaced"``; + - different bytes otherwise (someone edited the DDS copy): raise + :class:`DdsCopyConflict`, unless ``force=True``. + + The XML is re-checked first, so a hand-built :class:`SdiMetadata` + cannot put one record under another's UUID. + """ + iso.check(metadata.xml, metadata.uuid) + target = metadata_path(dds_path, metadata.uuid, folder=folder) + documents = self._documents_client() + + existing = documents.get(target) + if existing is None: + documents.put(target, metadata.xml, content_type=XML_CONTENT_TYPE) + return PushResult(target, "uploaded", metadata.uuid) + if existing == metadata.xml: + return PushResult(target, "unchanged", metadata.uuid) + if not force: + _refuse_unless_older(existing, metadata, target) + documents.put(target, metadata.xml, content_type=XML_CONTENT_TYPE, overwrite=True) + return PushResult(target, "replaced", metadata.uuid) + + def resolve_series(self, series_uuid: str) -> str: + """The UUID of the one release of ``series_uuid`` that is not superseded. + + Raises :class:`NotCurrentError`, listing every release, if there are + zero or several, rather than guessing by date. + """ + series_uuid = _require_uuid(series_uuid) + series = self.get_xml(series_uuid) + candidates = self._catalogue.search_by_uuid(list(series.record.children)) + current = [c for c in candidates if c.status != "superseded"] + if len(current) != 1: + raise NotCurrentError(series_uuid, candidates) + return current[0].uuid + + # -- internals -------------------------------------------------------- + + def _documents_client(self) -> DocumentsClient: + if self._documents is None: + self._documents = DocumentsClient.from_env() + return self._documents + + +def metadata_path(dds_path: str, uuid: str, *, folder: str = DEFAULT_FOLDER) -> str: + """``{dds_path}/{folder}/{uuid}.xml``, with slashes normalised.""" + parent = dds_path.strip().strip("/") + sub = folder.strip().strip("/").lower() + if not parent: + raise ValueError("dds_path must not be empty") + if not sub: + raise ValueError("folder must not be empty") + return f"{parent}/{sub}/{_require_uuid(uuid)}.xml" + + +def _require_uuid(uuid: str) -> str: + uuid = uuid.strip() + if not uuid: + raise ValueError("uuid must not be empty") + if "/" in uuid: + raise ValueError(f"invalid SDI UUID {uuid!r}") + return uuid + + +def _refuse_unless_older(existing: bytes, metadata: SdiMetadata, target: str) -> None: + """Raise :class:`DdsCopyConflict` unless the DDS copy is older than SDI's.""" + try: + theirs = iso.parse(existing).date_stamp + except NotIso19115_3: + raise DdsCopyConflict( + f"{target} exists but is not an ISO 19115-3 record; pass force=True to replace it" + ) from None + ours = metadata.record.date_stamp + if theirs is not None and ours is not None and ours > theirs: + return + raise DdsCopyConflict( + f"{target} differs from SDI record {metadata.uuid!r} and is not older " + f"(DDS {theirs.isoformat() if theirs else 'undated'}, " + f"SDI {ours.isoformat() if ours else 'undated'}); it may have been edited in DDS. " + "Pass force=True to replace it." + ) diff --git a/src/eea_datalakehouse/sdi/errors.py b/src/eea_datalakehouse/sdi/errors.py new file mode 100644 index 0000000..cf235f6 --- /dev/null +++ b/src/eea_datalakehouse/sdi/errors.py @@ -0,0 +1,78 @@ +"""Errors raised by the SDI controller.""" + +from __future__ import annotations + +from dataclasses import dataclass + + +class SdiError(RuntimeError): + """Base class for SDI controller errors.""" + + +class SdiApiError(SdiError): + """The SDI REST API returned a non-success status.""" + + def __init__(self, status_code: int, message: str, *, where: str) -> None: + super().__init__(f"SDI API error {status_code} on {where}: {message}") + self.status_code = status_code + self.message = message + self.where = where + + +class SdiNotFound(SdiApiError): + """No SDI record has that UUID (404).""" + + +class SdiAuthError(SdiApiError): + """SDI refused the request (401/403): the record is not public. + + Set ``SDI_USERNAME`` / ``SDI_PASSWORD`` for an account that may read it. + """ + + +class NotIso19115_3(SdiError): # noqa: N801 — the standard's name + """The document is not an ISO 19115-3 ``mdb:MD_Metadata`` record.""" + + +class UuidMismatch(SdiError): + """The XML's ``metadataIdentifier`` is not the UUID that was asked for.""" + + +class DdsCopyConflict(SdiError): + """The copy already in DDS differs and is not older than the SDI record. + + Someone probably edited it in DDS. Extract never silently overwrites that; + pass ``force=True`` to replace it anyway. + """ + + +@dataclass(frozen=True) +class SeriesCandidate: + """One release of a series, as SDI's search index describes it.""" + + uuid: str + title: str | None + edition: str | None + publication_date: str | None + status: str | None + + +class NotCurrentError(SdiError): + """A series has zero or several releases that are not ``superseded``. + + ``candidates`` lists every release so the caller can pick a UUID. + """ + + def __init__(self, series_uuid: str, candidates: list[SeriesCandidate]) -> None: + current = [c for c in candidates if c.status != "superseded"] + lines = "\n".join( + f" {c.uuid} {c.status or 'current'} ed. {c.edition or '-'} " + f"{c.publication_date or '-'} {c.title or ''}" + for c in candidates + ) + super().__init__( + f"series {series_uuid!r} has {len(current)} current release(s), expected 1; " + f"pass the release UUID explicitly:\n{lines}" + ) + self.series_uuid = series_uuid + self.candidates = candidates diff --git a/src/eea_datalakehouse/sdi/iso.py b/src/eea_datalakehouse/sdi/iso.py new file mode 100644 index 0000000..8009173 --- /dev/null +++ b/src/eea_datalakehouse/sdi/iso.py @@ -0,0 +1,133 @@ +"""Read the few ISO 19115-3 fields the SDI controller checks. + +Uses :mod:`xml.etree.ElementTree` with explicit namespaces. Python's expat +parser does not fetch external entities, so an SDI response cannot make it +read local files or the network. +""" + +from __future__ import annotations + +import xml.etree.ElementTree as ET +from dataclasses import dataclass +from datetime import UTC, date, datetime + +from .errors import NotIso19115_3, UuidMismatch + +NS = { + "mdb": "http://standards.iso.org/iso/19115/-3/mdb/2.0", + "mcc": "http://standards.iso.org/iso/19115/-3/mcc/1.0", + "mri": "http://standards.iso.org/iso/19115/-3/mri/1.0", + "cit": "http://standards.iso.org/iso/19115/-3/cit/2.0", + "gco": "http://standards.iso.org/iso/19115/-3/gco/1.0", + "gcx": "http://standards.iso.org/iso/19115/-3/gcx/1.0", +} +ROOT_TAG = f"{{{NS['mdb']}}}MD_Metadata" + +_IDENTIFIER = "mdb:metadataIdentifier/mcc:MD_Identifier/mcc:code/gco:CharacterString" +_SCOPE = "mdb:metadataScope/mdb:MD_MetadataScope/mdb:resourceScope/mcc:MD_ScopeCode" +_CITATION = "mdb:identificationInfo/*/mri:citation/cit:CI_Citation" +_DATE_INFO = "mdb:dateInfo/cit:CI_Date" +_ASSOCIATED = "mdb:identificationInfo/*/mri:associatedResource/mri:MD_AssociatedResource" + + +@dataclass(frozen=True) +class IsoRecord: + """What the controller needs to know about one ISO 19115-3 record.""" + + uuid: str + title: str | None + edition: str | None + hierarchy_level: str | None + created: datetime | None + revised: datetime | None + children: tuple[str, ...] = () + + @property + def date_stamp(self) -> datetime | None: + """The metadata's last change: its revision date, else its creation date.""" + return self.revised or self.created + + +def parse(xml: bytes) -> IsoRecord: + """Parse ``xml``; raise :class:`NotIso19115_3` unless it is ``mdb:MD_Metadata``.""" + try: + root = ET.fromstring(xml) + except ET.ParseError as exc: + raise NotIso19115_3(f"not well-formed XML: {exc}") from exc + if root.tag != ROOT_TAG: + raise NotIso19115_3(f"expected mdb:MD_Metadata (ISO 19115-3), got {root.tag!r}") + + citation = root.find(_CITATION, NS) + scope = root.find(_SCOPE, NS) + dates: dict[str, datetime] = {} + for ci_date in root.findall(_DATE_INFO, NS): + kind = ci_date.find("cit:dateType/cit:CI_DateTypeCode", NS) + when = _parse_date( + _text(ci_date, "cit:date/gco:DateTime") or _text(ci_date, "cit:date/gco:Date") + ) + if kind is not None and when is not None: + key = kind.get("codeListValue", "") + dates[key] = max(when, dates[key]) if key in dates else when + + children = tuple( + ref + for assoc in root.findall(_ASSOCIATED, NS) + if _code(assoc, "mri:associationType/mri:DS_AssociationTypeCode") == "isComposedOf" + and (ref := _attr(assoc, "mri:metadataReference", "uuidref")) + ) + + return IsoRecord( + uuid=_text(root, _IDENTIFIER) or "", + title=_text(citation, "cit:title") if citation is not None else None, + edition=_text(citation, "cit:edition") if citation is not None else None, + hierarchy_level=scope.get("codeListValue") if scope is not None else None, + created=dates.get("creation"), + revised=dates.get("revision"), + children=children, + ) + + +def check(xml: bytes, uuid: str) -> IsoRecord: + """Parse ``xml`` and raise :class:`UuidMismatch` unless it identifies ``uuid``.""" + record = parse(xml) + if record.uuid != uuid: + raise UuidMismatch(f"expected SDI record {uuid!r} but the XML identifies {record.uuid!r}") + return record + + +def _text(node: ET.Element, path: str) -> str | None: + """Text of a ``gco:CharacterString`` / ``gcx:Anchor`` (or a plain element).""" + el = node.find(path, NS) + if el is None: + return None + if el.text is None or not el.text.strip(): + inner = el.find("gco:CharacterString", NS) + if inner is None: + inner = el.find("gcx:Anchor", NS) + el = inner if inner is not None else el + text = (el.text or "").strip() + return text or None + + +def _code(node: ET.Element, path: str) -> str | None: + el = node.find(path, NS) + return el.get("codeListValue") if el is not None else None + + +def _attr(node: ET.Element, path: str, name: str) -> str | None: + el = node.find(path, NS) + return el.get(name) if el is not None else None + + +def _parse_date(value: str | None) -> datetime | None: + if not value: + return None + try: + parsed = datetime.fromisoformat(value) + except ValueError: + try: + parsed = datetime.combine(date.fromisoformat(value), datetime.min.time()) + except ValueError: + return None + # Dates without a zone are taken as UTC so every value compares. + return parsed if parsed.tzinfo else parsed.replace(tzinfo=UTC) diff --git a/src/eea_datalakehouse/sdi/session.py b/src/eea_datalakehouse/sdi/session.py new file mode 100644 index 0000000..ec44401 --- /dev/null +++ b/src/eea_datalakehouse/sdi/session.py @@ -0,0 +1,266 @@ +"""Notebook-facing sessions behind the `%sdi` and `%metadata` magics. + +Two sessions, one per magic, so reading SDI and writing DDS stay separate:: + + sdi = SdiSession() # %sdi — reads the SDI catalogue + sdi.get_xml("070d9baa-448d-4168-8514-7dadb3ad876d") + + metadata = MetadataSession(last=lambda: sdi.last) # %metadata — writes DDS + metadata.push_to_dds("catalog/water_management_resources/bathing_water/bwd") + +- **`SdiSession` remembers the last record** `get_xml` fetched, and + `MetadataSession.push_to_dds(dds_path)` uploads it by default, so a notebook + doesn't have to hold it in a variable. +- **Every failure is a `SdiSessionError` / `MetadataSessionError`**, so the + magic can print it as one short line instead of a traceback (the original + error is chained as ``__cause__``). + +Connection details come from the kernel environment, as `%catalog` / +`%ingest` read theirs: + +- `SdiSession`: ``SDI_API_URL`` (+ optional ``SDI_USERNAME`` / + ``SDI_PASSWORD``). It needs no DDS configuration at all. +- `MetadataSession`: ``DDS_BASE_URL`` and the Dremio identity. The identity is + ``_DREMIO_USER`` / ``_DREMIO_PWD`` (what `%ingest` reads), falling back to + ``DREMIO_USERNAME`` / ``DREMIO_TOKEN`` (what `%catalog` reads). + +**`DDS_BASE_URL` comes from a `.env` file first.** When a `MetadataSession` is +built, it looks for `.env` in the notebook's working directory, then in each +parent directory, and takes `DDS_BASE_URL` from the first one found. Only when +no `.env` sets it does the kernel environment's ``DDS_BASE_URL`` apply. +Nothing else is read from the file, and it never changes ``os.environ``. +`dds_base_url()` shows which URL is in use and where it came from. +""" + +from __future__ import annotations + +import os +from collections.abc import Callable, Mapping +from functools import wraps +from pathlib import Path +from typing import Any, TypeVar + +import httpx + +from ..dds_documents import DocumentsApiError, DocumentsClient +from ..dds_ingestion.credentials import ( + ENV_BASE_URL, + ENV_PWD, + ENV_USER, + DremioCreds, + MissingCredentialsError, + load_base_url, +) +from .catalogue import SdiCatalogue +from .controller import DEFAULT_FOLDER, PushResult, SdiController, SdiMetadata +from .errors import SdiError + +# What `%catalog` reads (see notebook/magics.py's _build_catalog_session). +ENV_CATALOG_USER = "DREMIO_USERNAME" +ENV_CATALOG_TOKEN = "DREMIO_TOKEN" # noqa: S105 — env var name, not a secret value +DOTENV_NAME = ".env" + +_F = TypeVar("_F", bound=Callable[..., Any]) + +# Failures a notebook user can act on; anything else is a bug and keeps its traceback. +_EXPECTED = (SdiError, DocumentsApiError, MissingCredentialsError, ValueError, httpx.HTTPError) + + +class SdiSessionError(RuntimeError): + """Any failure of an `SdiSession` call; the original error is ``__cause__``.""" + + +class MetadataSessionError(RuntimeError): + """Any failure of a `MetadataSession` call; the original error is ``__cause__``.""" + + +def _friendly(error: type[RuntimeError]) -> Callable[[_F], _F]: + """Re-raise every expected failure of the decorated method as ``error``.""" + + def decorate(method: _F) -> _F: + @wraps(method) + def wrapper(*args: Any, **kwargs: Any) -> Any: + try: + return method(*args, **kwargs) + except error: + raise + except _EXPECTED as exc: + raise error(str(exc) or type(exc).__name__) from exc + + return wrapper # type: ignore[return-value] + + return decorate + + +def load_dds_creds(env: Mapping[str, str] | None = None) -> DremioCreds: + """Dremio identity for the DDS upload: `%ingest`'s variables, else `%catalog`'s.""" + source = os.environ if env is None else env + for user_var, token_var in ((ENV_USER, ENV_PWD), (ENV_CATALOG_USER, ENV_CATALOG_TOKEN)): + user, token = source.get(user_var), source.get(token_var) + if user and token: + return DremioCreds(username=user, _password=token) + raise MissingCredentialsError( + f"pushing to DDS needs a Dremio identity in the kernel environment: " + f"{ENV_USER}/{ENV_PWD} or {ENV_CATALOG_USER}/{ENV_CATALOG_TOKEN}" + ) + + +def find_dotenv(start: str | Path | None = None) -> Path | None: + """The first ``.env`` in ``start`` (default: the working directory) or a parent.""" + here = Path.cwd() if start is None else Path(start) + for folder in (here, *here.resolve().parents): + candidate = folder / DOTENV_NAME + if candidate.is_file(): + return candidate + return None + + +def read_dotenv_value(path: Path, key: str) -> str | None: + """``key``'s value in the ``KEY=VALUE`` file ``path``, or ``None``. + + Blank lines, ``#`` comments and a leading ``export`` are allowed; one pair + of surrounding quotes is stripped. An empty value counts as unset. + """ + found: str | None = None + for raw in path.read_text(encoding="utf-8").splitlines(): + line = raw.strip() + if not line or line.startswith("#"): + continue + line = line.removeprefix("export ").strip() + name, sep, value = line.partition("=") + if not sep or name.strip() != key: + continue + value = value.strip() + if len(value) >= 2 and value[0] == value[-1] and value[0] in "\"'": + value = value[1:-1] + found = value or None # the last assignment wins, as in a shell + return found + + +class SdiSession: + """`get_xml` / `resolve_series` for a notebook — what `%sdi` dispatches onto. + + Reads the SDI catalogue only; pushing to DDS is `MetadataSession`'s job. + ``controller`` is built from the kernel environment if not given (see the + module docstring). One session lives for the whole kernel under `%sdi`. + """ + + def __init__( + self, + *, + controller: SdiController | None = None, + env: Mapping[str, str] | None = None, + ) -> None: + self._env = env + self._controller = controller + self._last: SdiMetadata | None = None + + def __repr__(self) -> str: + last = self._last.uuid if self._last else None + return f"SdiSession(last={last!r})" + + @property + def last(self) -> SdiMetadata | None: + """The record the last `get_xml` fetched, if any.""" + return self._last + + @_friendly(SdiSessionError) + def get_xml(self, uuid: str) -> SdiMetadata: + """Download the ISO 19115-3 XML of SDI record `uuid` and remember it.""" + self._last = self._sdi().get_xml(uuid) + return self._last + + @_friendly(SdiSessionError) + def resolve_series(self, series_uuid: str) -> str: + """The UUID of the one release of `series_uuid` that is not superseded.""" + return self._sdi().resolve_series(series_uuid) + + # -- internals -------------------------------------------------------- + + def _sdi(self) -> SdiController: + if self._controller is None: + env = dict(os.environ if self._env is None else self._env) + self._controller = SdiController(SdiCatalogue.from_env(env)) + return self._controller + + +class MetadataSession: + """`push_to_dds` / `dds_base_url` for a notebook — what `%metadata` dispatches onto. + + ``last`` returns the record to push by default — under the magics, the + `%sdi` session's last `get_xml` result. ``controller`` is built from the + kernel environment on the first push if not given; ``dotenv_path`` names + the ``.env`` to read ``DDS_BASE_URL`` from (by default searched for from + the working directory upwards, once, when the session is built). + """ + + def __init__( + self, + *, + last: Callable[[], SdiMetadata | None] | None = None, + controller: SdiController | None = None, + env: Mapping[str, str] | None = None, + dotenv_path: str | Path | None = None, + ) -> None: + self._last = last or (lambda: None) + self._controller = controller + self._env = env + self._dotenv_path = Path(dotenv_path) if dotenv_path is not None else find_dotenv() + self._dotenv_dds_url = ( + read_dotenv_value(self._dotenv_path, ENV_BASE_URL) + if self._dotenv_path is not None and self._dotenv_path.is_file() + else None + ) + + def __repr__(self) -> str: + return f"MetadataSession(dotenv={str(self._dotenv_path) if self._dotenv_path else None!r})" + + @_friendly(MetadataSessionError) + def push_to_dds( + self, + dds_path: str, + metadata: SdiMetadata | None = None, + folder: str = DEFAULT_FOLDER, + force: bool = False, + ) -> PushResult: + """Upload `metadata` (default: the last `%sdi get_xml` result) to + `{dds_path}/{folder}/{uuid}.xml` in DDS.""" + metadata = metadata or self._last() + if metadata is None: + raise MetadataSessionError("nothing to push yet — run %sdi get_xml(uuid) first") + return self._push_controller().push_to_dds( + metadata, dds_path, folder=folder, force=force + ) + + @_friendly(MetadataSessionError) + def dds_base_url(self) -> str: + """The DDS URL `push_to_dds` uses, and where it came from.""" + url, source = self._resolve_dds_base_url() + return f"{url} (from {source})" + + # -- internals -------------------------------------------------------- + + def _push_controller(self) -> SdiController: + if self._controller is None: + env = dict(os.environ if self._env is None else self._env) + # push_to_dds never calls SDI, so the catalogue is never used here. + self._controller = SdiController( + documents=DocumentsClient(self._resolve_dds_base_url()[0], load_dds_creds(env)) + ) + return self._controller + + def _resolve_dds_base_url(self) -> tuple[str, str]: + """``(url, source)``: the ``.env`` file's ``DDS_BASE_URL``, else the kernel's.""" + if self._dotenv_dds_url and self._dotenv_path is not None: + return self._dotenv_dds_url.rstrip("/"), str(self._dotenv_path) + env = dict(os.environ if self._env is None else self._env) + try: + return load_base_url(env), "the kernel environment" + except MissingCredentialsError: + where = ( + f"{self._dotenv_path} has none and " if self._dotenv_path else "no .env found and " + ) + raise MissingCredentialsError( + f"pushing to DDS needs {ENV_BASE_URL}: {where}it is not set in the kernel " + "environment" + ) from None diff --git a/tests/dds_documents/__init__.py b/tests/dds_documents/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/tests/dds_documents/test_client.py b/tests/dds_documents/test_client.py new file mode 100644 index 0000000..e36714a --- /dev/null +++ b/tests/dds_documents/test_client.py @@ -0,0 +1,104 @@ +from __future__ import annotations + +import httpx +import pytest +import respx + +from eea_datalakehouse.dds_documents import ( + DocumentExistsError, + DocumentsApiError, + DocumentsClient, +) +from eea_datalakehouse.dds_ingestion.credentials import DremioCreds + +BASE_URL = "https://dds.example.test" + + +@pytest.fixture +def client() -> DocumentsClient: + return DocumentsClient(BASE_URL, DremioCreds(username="alice", _password="s3cr3t-pwd")) + + +@respx.mock +def test_put_sends_body_type_and_bearer(client: DocumentsClient) -> None: + route = respx.put(f"{BASE_URL}/api/v1/files/a/b/x.xml").mock( + return_value=httpx.Response(201) + ) + + client.put("/a/b/x.xml", b"", content_type="application/xml") + + request = route.calls.last.request + assert request.content == b"" + assert request.headers["Content-Type"] == "application/xml" + assert request.headers["Authorization"] == "Bearer s3cr3t-pwd" + assert request.url.params["overwrite"] == "false" + + +@respx.mock +def test_put_encodes_spaces_in_path(client: DocumentsClient) -> None: + route = respx.put(f"{BASE_URL}/api/v1/files/a%20b/x.xml").mock( + return_value=httpx.Response(201) + ) + + client.put("a b/x.xml", b"") + + assert route.called + + +@respx.mock +def test_put_conflict_is_document_exists(client: DocumentsClient) -> None: + respx.put(f"{BASE_URL}/api/v1/files/a/x.xml").mock( + return_value=httpx.Response(409, json={"error": "exists", "message": "file exists"}) + ) + + with pytest.raises(DocumentExistsError, match="file exists"): + client.put("a/x.xml", b"") + + +@respx.mock +def test_put_error_names_the_call(client: DocumentsClient) -> None: + respx.put(f"{BASE_URL}/api/v1/files/a/x.xml").mock( + return_value=httpx.Response(403, json={"error": "forbidden", "message": "no write access"}) + ) + + with pytest.raises(DocumentsApiError, match=r"403 on PUT /api/v1/files/a/x.xml: no write"): + client.put("a/x.xml", b"") + + +@respx.mock +def test_get_and_exists(client: DocumentsClient) -> None: + respx.get(f"{BASE_URL}/api/v1/files/a/x.xml").mock( + return_value=httpx.Response(200, content=b"") + ) + respx.get(f"{BASE_URL}/api/v1/files/a/missing.xml").mock(return_value=httpx.Response(404)) + + assert client.get("a/x.xml") == b"" + assert client.exists("a/x.xml") + assert client.get("a/missing.xml") is None + assert not client.exists("a/missing.xml") + + +@respx.mock +@pytest.mark.parametrize( + "payload", [{"files": [{"path": "a/x.xml"}, {"path": "a/y.xml"}]}, ["a/x.xml", "a/y.xml"]] +) +def test_list(client: DocumentsClient, payload: object) -> None: + route = respx.get(f"{BASE_URL}/api/v1/files").mock( + return_value=httpx.Response(200, json=payload) + ) + + assert client.list("/a/") == ["a/x.xml", "a/y.xml"] + assert route.calls.last.request.url.params["prefix"] == "a" + + +@pytest.mark.parametrize("bad", ["", "/", "a/../b", "a//b"]) +def test_rejects_bad_paths(client: DocumentsClient, bad: str) -> None: + with pytest.raises(ValueError): + client.get(bad) + + +def test_from_env_and_repr_hide_the_token() -> None: + client = DocumentsClient.from_env( + {"DDS_BASE_URL": f"{BASE_URL}/", "_DREMIO_USER": "u", "_DREMIO_PWD": "s3cr3t"} + ) + assert repr(client) == f"DocumentsClient(base_url={BASE_URL!r})" diff --git a/tests/notebook/test_magics.py b/tests/notebook/test_magics.py index 6e4dbb3..4f5aa89 100644 --- a/tests/notebook/test_magics.py +++ b/tests/notebook/test_magics.py @@ -85,6 +85,8 @@ def ip(shell: Any) -> Any: instance = _magics_instance(shell) instance._catalog_session = None instance._ingest_session = None + instance._sdi_session = None + instance._metadata_session = None return shell @@ -399,3 +401,115 @@ def test_ingest_help_lists_methods( assert "retry=False, max_retries=3" in table_html assert "str | None" not in table_html # no Python type-hint syntax leaking into the table assert _magics_instance(ip)._ingest_session is None + + +def test_sdi_magic_builds_one_session_and_dispatches_onto_it( + ip: Any, monkeypatch: pytest.MonkeyPatch +) -> None: + calls: list[str] = [] + + class _FakeSdiSession: + def get_xml(self, uuid: str) -> str: + calls.append(uuid) + return f"metadata {uuid}" + + monkeypatch.setattr(magics_module, "_build_sdi_session", _FakeSdiSession) + + assert ip.run_line_magic("sdi", 'get_xml("u-1")') == "metadata u-1" + first = _magics_instance(ip)._sdi_session + ip.run_line_magic("sdi", 'get_xml("u-2")') + + assert calls == ["u-1", "u-2"] + assert _magics_instance(ip)._sdi_session is first + + +def test_sdi_session_error_is_printed_not_raised( + ip: Any, monkeypatch: pytest.MonkeyPatch, capsys: pytest.CaptureFixture[str] +) -> None: + monkeypatch.setenv("SDI_API_URL", "https://sdi.example.test/catalogue") + + assert ip.run_line_magic("sdi", 'get_xml("")') is None + + assert "sdi error: uuid must not be empty" in capsys.readouterr().out + + +def test_metadata_push_before_sdi_get_xml_is_printed_not_raised( + ip: Any, capsys: pytest.CaptureFixture[str] +) -> None: + assert ip.run_line_magic("metadata", 'push_to_dds("a/b")') is None + + assert "metadata error: nothing to push yet — run %sdi get_xml" in capsys.readouterr().out + + +def test_metadata_pushes_the_record_sdi_fetched_last( + ip: Any, monkeypatch: pytest.MonkeyPatch +) -> None: + pushed: list[object] = [] + + class _FakeSdiSession: + last: object = None + + def get_xml(self, uuid: str) -> str: + self.last = f"record {uuid}" + return self.last + + class _FakeMetadataSession: + def __init__(self, last: Any) -> None: + self._last = last + + def push_to_dds(self, dds_path: str) -> str: + pushed.append(self._last()) + return f"pushed to {dds_path}" + + monkeypatch.setattr(magics_module, "_build_sdi_session", _FakeSdiSession) + monkeypatch.setattr(magics_module, "_build_metadata_session", _FakeMetadataSession) + + # %metadata built first: the record is still looked up at push time. + ip.run_line_magic("metadata", "help") + assert ip.run_line_magic("metadata", 'push_to_dds("a/b")') == "pushed to a/b" + ip.run_line_magic("sdi", 'get_xml("u-1")') + ip.run_line_magic("metadata", 'push_to_dds("a/b")') + + assert pushed == [None, "record u-1"] + + +def test_sdi_empty_line_prints_usage(ip: Any, capsys: pytest.CaptureFixture[str]) -> None: + ip.run_line_magic("sdi", "") + + assert "usage: %sdi get_xml" in capsys.readouterr().out + + +@pytest.mark.parametrize("line", ["help", "help()"]) +def test_sdi_help_lists_methods_without_building_a_session( + ip: Any, monkeypatch: pytest.MonkeyPatch, capsys: pytest.CaptureFixture[str], line: str +) -> None: + displayed = [] + monkeypatch.setattr(magics_module, "display", displayed.append) + + ip.run_line_magic("sdi", line) + + out = capsys.readouterr().out + assert "%sdi methods" in out + assert "SDI_API_URL" in out + table_html = displayed[0].data + assert ">get_xml<" in table_html and ">resolve_series<" in table_html + assert ">push_to_dds<" not in table_html # moved to %metadata + assert _magics_instance(ip)._sdi_session is None + + +@pytest.mark.parametrize("line", ["help", "help()"]) +def test_metadata_help_lists_methods_without_building_a_session( + ip: Any, monkeypatch: pytest.MonkeyPatch, capsys: pytest.CaptureFixture[str], line: str +) -> None: + displayed = [] + monkeypatch.setattr(magics_module, "display", displayed.append) + + ip.run_line_magic("metadata", line) + + out = capsys.readouterr().out + assert "%metadata methods" in out + assert "DDS_BASE_URL" in out + table_html = displayed[0].data + assert ">push_to_dds<" in table_html and ">dds_base_url<" in table_html + assert "dds_path, metadata=None, folder='metadata', force=False" in table_html + assert _magics_instance(ip)._metadata_session is None diff --git a/tests/sdi/__init__.py b/tests/sdi/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/tests/sdi/conftest.py b/tests/sdi/conftest.py new file mode 100644 index 0000000..fff71b9 --- /dev/null +++ b/tests/sdi/conftest.py @@ -0,0 +1,104 @@ +from __future__ import annotations + +import httpx +import pytest + +from eea_datalakehouse.dds_documents import DocumentsClient +from eea_datalakehouse.dds_ingestion.credentials import DremioCreds +from eea_datalakehouse.sdi import SdiCatalogue, SdiConfig, SdiController + +SDI_URL = "https://sdi.example.test/catalogue" +API = f"{SDI_URL}/srv/api" +DDS_URL = "https://dds.example.test" +UUID = "070d9baa-448d-4168-8514-7dadb3ad876d" +SERIES_UUID = "c3858959-90da-4c1b-b9ca-492db0e514df" + +_NAMESPACES = ( + 'xmlns:mdb="http://standards.iso.org/iso/19115/-3/mdb/2.0" ' + 'xmlns:mcc="http://standards.iso.org/iso/19115/-3/mcc/1.0" ' + 'xmlns:mri="http://standards.iso.org/iso/19115/-3/mri/1.0" ' + 'xmlns:cit="http://standards.iso.org/iso/19115/-3/cit/2.0" ' + 'xmlns:gco="http://standards.iso.org/iso/19115/-3/gco/1.0"' +) + + +def iso_xml( + uuid: str = UUID, + *, + revised: str = "2026-09-25T07:36:46.809835Z", + title: str = "Bathing Water Directive - Status of bathing water, 2025 v.1.0", + children: tuple[str, ...] = (), + scope: str = "nonGeographicDataset", +) -> bytes: + """A trimmed ISO 19115-3 record with the same structure SDI serves.""" + assoc = "".join( + f""" + + + + + + """ + for child in children + ) + return f""" + + + {uuid} + + + + + + 2026-06-01T07:54:39.676514Z + + + + {revised} + + + + + {title} + 01.00 + {assoc} + + +""".encode() + + +def search_hit(uuid: str, *, superseded: bool, edition: str = "01.00") -> dict[str, object]: + source: dict[str, object] = { + "uuid": uuid, + "resourceTitleObject": {"default": f"release {uuid}"}, + "resourceEdition": edition, + "publicationDateForResource": ["2026-06-02"], + } + if superseded: + source["cl_status"] = [{"key": "superseded", "default": "Superseded"}] + return {"_source": source} + + +def record_url(uuid: str = UUID) -> str: + return f"{API}/records/{uuid}/formatters/xml" + + +def dds_file_url(path: str) -> str: + return f"{DDS_URL}/api/v1/files/{path}" + + +@pytest.fixture +def creds() -> DremioCreds: + return DremioCreds(username="alice", _password="s3cr3t-pwd") + + +@pytest.fixture +def sdi(creds: DremioCreds) -> SdiController: + return SdiController( + SdiCatalogue(SdiConfig(api_url=SDI_URL)), + DocumentsClient(DDS_URL, creds), + ) + + +def not_found() -> httpx.Response: + return httpx.Response(404, json={"error": "not_found", "message": "no such file"}) diff --git a/tests/sdi/test_catalogue.py b/tests/sdi/test_catalogue.py new file mode 100644 index 0000000..f55bf96 --- /dev/null +++ b/tests/sdi/test_catalogue.py @@ -0,0 +1,110 @@ +from __future__ import annotations + +import json +import os + +import httpx +import pytest +import respx + +from eea_datalakehouse.sdi import ( + SdiApiError, + SdiAuthError, + SdiCatalogue, + SdiConfig, + SdiController, + SdiNotFound, +) + +from .conftest import API, SDI_URL, UUID, iso_xml, record_url, search_hit + + +@pytest.fixture +def catalogue() -> SdiCatalogue: + return SdiCatalogue(SdiConfig(api_url=SDI_URL)) + + +@respx.mock +def test_get_xml_returns_bytes_unchanged(catalogue: SdiCatalogue) -> None: + route = respx.get(record_url()).mock(return_value=httpx.Response(200, content=iso_xml())) + + assert catalogue.get_xml(UUID) == iso_xml() + assert route.calls.last.request.headers["Accept"] == "application/xml" + assert "Authorization" not in route.calls.last.request.headers + + +@respx.mock +def test_get_xml_sends_basic_auth_when_configured() -> None: + route = respx.get(record_url()).mock(return_value=httpx.Response(200, content=iso_xml())) + catalogue = SdiCatalogue(SdiConfig(api_url=SDI_URL, username="u", _password="p")) + + catalogue.get_xml(UUID) + + assert route.calls.last.request.headers["Authorization"].startswith("Basic ") + + +@respx.mock +@pytest.mark.parametrize( + ("status", "error"), [(404, SdiNotFound), (401, SdiAuthError), (403, SdiAuthError)] +) +def test_get_xml_errors(catalogue: SdiCatalogue, status: int, error: type[Exception]) -> None: + respx.get(record_url()).mock(return_value=httpx.Response(status)) + + with pytest.raises(error): + catalogue.get_xml(UUID) + + +@respx.mock +def test_get_xml_other_error_names_the_call(catalogue: SdiCatalogue) -> None: + respx.get(record_url()).mock(return_value=httpx.Response(502, text="Bad Gateway")) + + with pytest.raises(SdiApiError, match=r"502 on GET /srv/api/records/.*Bad Gateway"): + catalogue.get_xml(UUID) + + +@respx.mock +def test_search_by_uuid(catalogue: SdiCatalogue) -> None: + route = respx.post(f"{API}/search/records/_search").mock( + return_value=httpx.Response( + 200, + json={ + "hits": { + "hits": [search_hit("a", superseded=True), search_hit("b", superseded=False)] + } + }, + ) + ) + + hits = catalogue.search_by_uuid(["a", "b"]) + + assert [(h.uuid, h.status) for h in hits] == [("a", "superseded"), ("b", None)] + assert hits[1].publication_date == "2026-06-02" + assert hits[1].title == "release b" + body = json.loads(route.calls.last.request.content) + assert body["query"] == {"terms": {"uuid": ["a", "b"]}} + assert body["size"] == 2 + + +def test_search_by_uuid_with_nothing_makes_no_call(catalogue: SdiCatalogue) -> None: + assert catalogue.search_by_uuid([]) == [] + + +def test_config_from_env_defaults_to_public_catalogue() -> None: + assert SdiConfig.from_env({}).api_url == "https://sdi.eea.europa.eu/catalogue" + assert SdiConfig.from_env({"SDI_API_URL": f"{SDI_URL}/"}).api_url == SDI_URL + + +def test_config_repr_hides_password() -> None: + config = SdiConfig.from_env({"SDI_USERNAME": "u", "SDI_PASSWORD": "hunter2"}) + assert "hunter2" not in repr(config) + assert config.auth is not None + + +@pytest.mark.skipif( + not os.environ.get("SDI_NETWORK_TESTS"), reason="set SDI_NETWORK_TESTS=1 to call live SDI" +) +def test_live_catalogue_extracts_bathing_water_release() -> None: + with SdiController.from_env() as sdi: + metadata = sdi.get_xml(UUID) + assert metadata.record.uuid == UUID + assert metadata.record.title diff --git a/tests/sdi/test_controller.py b/tests/sdi/test_controller.py new file mode 100644 index 0000000..5352d89 --- /dev/null +++ b/tests/sdi/test_controller.py @@ -0,0 +1,218 @@ +from __future__ import annotations + +import httpx +import pytest +import respx + +from eea_datalakehouse.sdi import ( + DdsCopyConflict, + NotCurrentError, + NotIso19115_3, + SdiController, + SdiMetadata, + SdiNotFound, + UuidMismatch, + metadata_path, +) + +from .conftest import ( + API, + SERIES_UUID, + UUID, + dds_file_url, + iso_xml, + not_found, + record_url, + search_hit, +) + +TARGET = f"water/bathing_water/metadata/{UUID}.xml" + + +def extracted(sdi: SdiController, xml: bytes | None = None) -> SdiMetadata: + with respx.mock: + respx.get(record_url()).mock(return_value=httpx.Response(200, content=xml or iso_xml())) + return sdi.get_xml(UUID) + + +# -- get_xml ------------------------------------------------------------ + + +def test_extract_returns_sdi_bytes_and_summary(sdi: SdiController) -> None: + metadata = extracted(sdi) + + assert metadata.uuid == UUID + assert metadata.xml == iso_xml() + assert metadata.record.edition == "01.00" + assert "bytes" in repr(metadata) + + +@respx.mock +def test_extract_404_is_not_found(sdi: SdiController) -> None: + respx.get(record_url()).mock(return_value=httpx.Response(404)) + + with pytest.raises(SdiNotFound): + sdi.get_xml(UUID) + + +@respx.mock +def test_extract_rejects_non_iso_19115_3(sdi: SdiController) -> None: + respx.get(record_url()).mock(return_value=httpx.Response(200, text="")) + + with pytest.raises(NotIso19115_3, match=UUID): + sdi.get_xml(UUID) + + +@respx.mock +def test_extract_rejects_another_records_xml(sdi: SdiController) -> None: + respx.get(record_url()).mock(return_value=httpx.Response(200, content=iso_xml("other"))) + + with pytest.raises(UuidMismatch): + sdi.get_xml(UUID) + + +@pytest.mark.parametrize("bad", ["", " ", "a/b"]) +def test_extract_rejects_bad_uuid(sdi: SdiController, bad: str) -> None: + with pytest.raises(ValueError): + sdi.get_xml(bad) + + +def test_provenance_tags_match_setmeta2wiki_shape(sdi: SdiController) -> None: + tags = {t["tag_name"]: t["tag_value"] for t in extracted(sdi).provenance_tags()} + + assert tags == { + "sdi_record_uuid": UUID, + "sdi_edition": "01.00", + "sdi_date_stamp": "2026-09-25T07:36:46.809835+00:00", + } + + +# -- push_to_dds ------------------------------------------------------------ + + +def test_metadata_path() -> None: + assert metadata_path("/water/bathing_water/", UUID) == TARGET + assert metadata_path("water", UUID, folder="/DOCS/") == f"water/docs/{UUID}.xml" + with pytest.raises(ValueError): + metadata_path("/", UUID) + + +def test_push_uploads_new_file(sdi: SdiController) -> None: + metadata = extracted(sdi) + with respx.mock: + respx.get(dds_file_url(TARGET)).mock(return_value=not_found()) + put = respx.put(dds_file_url(TARGET)).mock(return_value=httpx.Response(201)) + + result = sdi.push_to_dds(metadata, "/water/bathing_water/") + + assert (result.dds_path, result.action) == (TARGET, "uploaded") + request = put.calls.last.request + assert request.content == iso_xml() + assert request.headers["Authorization"] == "Bearer s3cr3t-pwd" + assert request.headers["Content-Type"] == "application/xml" + assert request.url.params["overwrite"] == "false" + + +def test_push_same_bytes_is_unchanged(sdi: SdiController) -> None: + metadata = extracted(sdi) + with respx.mock: + respx.get(dds_file_url(TARGET)).mock(return_value=httpx.Response(200, content=iso_xml())) + put = respx.put(dds_file_url(TARGET)) + + result = sdi.push_to_dds(metadata, "water/bathing_water") + + assert result.action == "unchanged" + assert not put.called + + +def test_push_replaces_older_copy(sdi: SdiController) -> None: + metadata = extracted(sdi) + older = iso_xml(revised="2026-06-02T00:00:00Z") + with respx.mock: + respx.get(dds_file_url(TARGET)).mock(return_value=httpx.Response(200, content=older)) + put = respx.put(dds_file_url(TARGET)).mock(return_value=httpx.Response(200)) + + result = sdi.push_to_dds(metadata, "water/bathing_water") + + assert result.action == "replaced" + assert put.calls.last.request.url.params["overwrite"] == "true" + + +@pytest.mark.parametrize( + "dds_copy", + [ + iso_xml(revised="2026-09-30T00:00:00Z"), # newer than SDI's + iso_xml(title="edited in DDS"), # same date, different content + b"not xml at all", + ], +) +def test_push_refuses_to_overwrite_dds_edits(sdi: SdiController, dds_copy: bytes) -> None: + metadata = extracted(sdi) + with respx.mock: + respx.get(dds_file_url(TARGET)).mock(return_value=httpx.Response(200, content=dds_copy)) + put = respx.put(dds_file_url(TARGET)) + + with pytest.raises(DdsCopyConflict, match="force=True"): + sdi.push_to_dds(metadata, "water/bathing_water") + + assert not put.called + + +def test_push_force_overwrites_dds_edits(sdi: SdiController) -> None: + metadata = extracted(sdi) + newer = iso_xml(revised="2026-09-30T00:00:00Z") + with respx.mock: + respx.get(dds_file_url(TARGET)).mock(return_value=httpx.Response(200, content=newer)) + respx.put(dds_file_url(TARGET)).mock(return_value=httpx.Response(200)) + + assert sdi.push_to_dds(metadata, "water/bathing_water", force=True).action == "replaced" + + +def test_push_refuses_mismatched_record(sdi: SdiController) -> None: + metadata = extracted(sdi) + forged = SdiMetadata(UUID, iso_xml("other"), metadata.record) + + with pytest.raises(UuidMismatch): + sdi.push_to_dds(forged, "water") + + +# -- resolve_series --------------------------------------------------------- + + +def _mock_series(hits: list[dict[str, object]]) -> None: + respx.get(record_url(SERIES_UUID)).mock( + return_value=httpx.Response( + 200, content=iso_xml(SERIES_UUID, scope="series", children=("old", UUID)) + ) + ) + respx.post(f"{API}/search/records/_search").mock( + return_value=httpx.Response(200, json={"hits": {"hits": hits}}) + ) + + +@respx.mock +def test_resolve_series_returns_the_one_current_release(sdi: SdiController) -> None: + _mock_series([search_hit("old", superseded=True), search_hit(UUID, superseded=False)]) + + assert sdi.resolve_series(SERIES_UUID) == UUID + + +@respx.mock +@pytest.mark.parametrize("current", [0, 2]) +def test_resolve_series_refuses_to_guess(sdi: SdiController, current: int) -> None: + _mock_series( + [search_hit("old", superseded=current == 0), search_hit(UUID, superseded=current == 0)] + ) + + with pytest.raises(NotCurrentError) as info: + sdi.resolve_series(SERIES_UUID) + + assert {c.uuid for c in info.value.candidates} == {"old", UUID} + assert f"has {current} current release(s)" in str(info.value) + + +# -- configuration ---------------------------------------------------------- + + +def test_repr_hides_credentials(sdi: SdiController) -> None: + assert "s3cr3t" not in repr(sdi) diff --git a/tests/sdi/test_iso.py b/tests/sdi/test_iso.py new file mode 100644 index 0000000..323968b --- /dev/null +++ b/tests/sdi/test_iso.py @@ -0,0 +1,54 @@ +from __future__ import annotations + +from datetime import UTC, datetime + +import pytest + +from eea_datalakehouse.sdi import NotIso19115_3, UuidMismatch, iso + +from .conftest import SERIES_UUID, UUID, iso_xml + + +def test_parse_reads_identifier_citation_scope_and_dates() -> None: + record = iso.parse(iso_xml()) + + assert record.uuid == UUID + assert record.title == "Bathing Water Directive - Status of bathing water, 2025 v.1.0" + assert record.edition == "01.00" + assert record.hierarchy_level == "nonGeographicDataset" + assert record.created == datetime(2026, 6, 1, 7, 54, 39, 676514, tzinfo=UTC) + assert record.date_stamp == datetime(2026, 9, 25, 7, 36, 46, 809835, tzinfo=UTC) + assert record.children == () + + +def test_parse_collects_series_children() -> None: + record = iso.parse(iso_xml(SERIES_UUID, scope="series", children=("a", "b"))) + + assert record.children == ("a", "b") + assert record.hierarchy_level == "series" + + +def test_plain_date_is_taken_as_utc() -> None: + xml = iso_xml().replace( + b"2026-09-25T07:36:46.809835Z", + b"2026-09-25", + ) + assert iso.parse(xml).revised == datetime(2026, 9, 25, tzinfo=UTC) + + +@pytest.mark.parametrize( + "xml", + [ + b'', + b"Service unavailable", + b"", + ], +) +def test_parse_rejects_anything_but_iso_19115_3(xml: bytes) -> None: + with pytest.raises(NotIso19115_3): + iso.parse(xml) + + +def test_check_rejects_another_records_xml() -> None: + with pytest.raises(UuidMismatch, match="identifies 'other'"): + iso.check(iso_xml("other"), UUID) diff --git a/tests/sdi/test_session.py b/tests/sdi/test_session.py new file mode 100644 index 0000000..7f3ea57 --- /dev/null +++ b/tests/sdi/test_session.py @@ -0,0 +1,191 @@ +from __future__ import annotations + +from pathlib import Path + +import httpx +import pytest +import respx + +from eea_datalakehouse.dds_ingestion.credentials import MissingCredentialsError +from eea_datalakehouse.sdi import ( + MetadataSession, + MetadataSessionError, + NotIso19115_3, + SdiController, + SdiSession, + SdiSessionError, +) +from eea_datalakehouse.sdi.session import find_dotenv, load_dds_creds, read_dotenv_value + +from .conftest import ( + DDS_URL, + SDI_URL, + UUID, + dds_file_url, + iso_xml, + not_found, + record_url, +) + +TARGET = f"water/metadata/{UUID}.xml" +ENV = {"SDI_API_URL": SDI_URL, "DDS_BASE_URL": DDS_URL} +NO_DOTENV = Path("/nonexistent/.env") + + +def _metadata_session(sdi: SdiSession | None = None, **kwargs: object) -> MetadataSession: + kwargs.setdefault("dotenv_path", NO_DOTENV) + kwargs.setdefault("env", ENV) + last = (lambda: sdi.last) if sdi is not None else None + return MetadataSession(last=last, **kwargs) # type: ignore[arg-type] + + +# -- SdiSession --------------------------------------------------------------- + + +@respx.mock +def test_get_xml_remembers_the_record_and_needs_no_dds_settings() -> None: + respx.get(record_url()).mock(return_value=httpx.Response(200, content=iso_xml())) + session = SdiSession(env={"SDI_API_URL": SDI_URL}) + + metadata = session.get_xml(UUID) + + assert session.last is metadata + assert metadata.uuid == UUID + assert repr(session) == f"SdiSession(last={UUID!r})" + + +@respx.mock +def test_sdi_errors_become_session_errors() -> None: + respx.get(record_url()).mock(return_value=httpx.Response(200, text="")) + + with pytest.raises(SdiSessionError) as info: + SdiSession(env=ENV).get_xml(UUID) + assert isinstance(info.value.__cause__, NotIso19115_3) + + +def test_sdi_session_has_no_push() -> None: + assert not hasattr(SdiSession, "push_to_dds") + assert not hasattr(SdiSession, "dds_base_url") + + +# -- MetadataSession ---------------------------------------------------------- + + +@respx.mock +def test_push_uploads_the_last_sdi_record_with_the_kernel_identity() -> None: + respx.get(record_url()).mock(return_value=httpx.Response(200, content=iso_xml())) + respx.get(dds_file_url(TARGET)).mock(return_value=not_found()) + put = respx.put(dds_file_url(TARGET)).mock(return_value=httpx.Response(201)) + sdi = SdiSession(env=ENV) + metadata = _metadata_session( + sdi, env={**ENV, "DREMIO_USERNAME": "carol", "DREMIO_TOKEN": "pat-c"} + ) + + sdi.get_xml(UUID) + result = metadata.push_to_dds("water") + + assert (result.action, result.dds_path) == ("uploaded", TARGET) + assert put.calls.last.request.headers["Authorization"] == "Bearer pat-c" + + +@respx.mock +def test_push_takes_an_explicit_record(sdi: SdiController) -> None: + respx.get(record_url()).mock(return_value=httpx.Response(200, content=iso_xml())) + respx.get(dds_file_url(TARGET)).mock(return_value=not_found()) + respx.put(dds_file_url(TARGET)).mock(return_value=httpx.Response(201)) + record = sdi.get_xml(UUID) + + result = _metadata_session(controller=sdi).push_to_dds("water", metadata=record) + + assert result.action == "uploaded" + + +def test_push_before_get_xml_is_a_session_error() -> None: + with pytest.raises(MetadataSessionError, match=r"run %sdi get_xml"): + _metadata_session(SdiSession(env=ENV)).push_to_dds("water") + + +@respx.mock +def test_push_without_dremio_identity_is_a_session_error() -> None: + respx.get(record_url()).mock(return_value=httpx.Response(200, content=iso_xml())) + sdi = SdiSession(env=ENV) + sdi.get_xml(UUID) + + with pytest.raises(MetadataSessionError, match="DREMIO_TOKEN") as info: + _metadata_session(sdi).push_to_dds("water") + assert isinstance(info.value.__cause__, MissingCredentialsError) + + +def test_load_dds_creds_prefers_ingest_variables() -> None: + env = {"_DREMIO_USER": "a", "_DREMIO_PWD": "1", "DREMIO_USERNAME": "b", "DREMIO_TOKEN": "2"} + assert load_dds_creds(env).username == "a" + assert load_dds_creds({"DREMIO_USERNAME": "b", "DREMIO_TOKEN": "2"}).username == "b" + with pytest.raises(MissingCredentialsError): + load_dds_creds({}) + + +# -- DDS_BASE_URL from .env ------------------------------------------------- + + +def test_dotenv_dds_url_wins_over_kernel_env(tmp_path: Path) -> None: + dotenv = tmp_path / ".env" + dotenv.write_text( + "# comment\nexport DDS_BASE_URL='https://dds.from-dotenv.test/'\nOTHER=x\n" + ) + session = _metadata_session(dotenv_path=dotenv) + + assert session.dds_base_url() == f"https://dds.from-dotenv.test (from {dotenv})" + + +@respx.mock +def test_push_goes_to_the_dotenv_url(tmp_path: Path) -> None: + (tmp_path / ".env").write_text("DDS_BASE_URL=https://dds.from-dotenv.test\n") + url = f"https://dds.from-dotenv.test/api/v1/files/{TARGET}" + respx.get(record_url()).mock(return_value=httpx.Response(200, content=iso_xml())) + respx.get(url).mock(return_value=not_found()) + put = respx.put(url).mock(return_value=httpx.Response(201)) + sdi = SdiSession(env=ENV) + sdi.get_xml(UUID) + + _metadata_session( + sdi, + dotenv_path=tmp_path / ".env", + env={**ENV, "_DREMIO_USER": "a", "_DREMIO_PWD": "p"}, + ).push_to_dds("water") + + assert put.called + + +def test_empty_dotenv_value_falls_back_to_kernel_env(tmp_path: Path) -> None: + (tmp_path / ".env").write_text("DDS_BASE_URL=\n") + session = _metadata_session(dotenv_path=tmp_path / ".env") + + assert session.dds_base_url() == f"{DDS_URL} (from the kernel environment)" + + +def test_no_dds_url_anywhere_names_both_places(tmp_path: Path) -> None: + (tmp_path / ".env").write_text("OTHER=x\n") + session = _metadata_session(dotenv_path=tmp_path / ".env", env={}) + + with pytest.raises(MetadataSessionError, match=r"\.env has none and it is not set"): + session.dds_base_url() + + +def test_find_dotenv_walks_up_from_the_working_directory( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch +) -> None: + (tmp_path / ".env").write_text("DDS_BASE_URL=https://dds.parent.test\n") + nested = tmp_path / "a" / "b" + nested.mkdir(parents=True) + monkeypatch.chdir(nested) + + assert find_dotenv() == tmp_path / ".env" + assert MetadataSession(env={}).dds_base_url().startswith("https://dds.parent.test") + + +def test_read_dotenv_value_last_assignment_wins(tmp_path: Path) -> None: + dotenv = tmp_path / ".env" + dotenv.write_text('DDS_BASE_URL=one\nDDS_BASE_URL="two"\n') + + assert read_dotenv_value(dotenv, "DDS_BASE_URL") == "two" + assert read_dotenv_value(dotenv, "MISSING") is None