diff --git a/examples/measurement/README.md b/examples/measurement/README.md new file mode 100644 index 000000000..6c780c97f --- /dev/null +++ b/examples/measurement/README.md @@ -0,0 +1,79 @@ +# Uploading a measurement run + +Data a lab delivers becomes Samples, Measurements and Properties on the platform in two steps that stay apart: +**parse** reads one lab's delivery and writes a *run document*; **upload** takes run documents and sends them. +Nothing in the uploader knows what an instrument is, and no parser talks to the platform. + +``` + delivery ──► parse_utk.py ──► run document ──► upload_run.py ──► platform + parse_nlr.py (parsed/*.json) +``` + +## Install + +```bash +python -m venv venv && . venv/bin/activate # Windows: venv\Scripts\activate +pip install -r requirements.txt +``` + +api-client, standata and esse are pinned to the branch carrying the REST endpoints, the instrument registry +entries and the Sample and Measurement schemas. They become version pins when it releases. + +## Credentials + +The uploader reads them from the environment. Either an API token from the platform's Preferences page: + +```bash +export ACCOUNT_ID=... AUTH_TOKEN=... +export MAT3RA_HOST=alphafilm.mat3ra.com # or pass --host +``` + +or, if you signed in through the browser, `OIDC_ACCESS_TOKEN`. + +## Parse + +Each parser reads one lab's delivery. Point it at the folder as delivered and give it the identifier written +on the physical piece — every Sample carries it, and it is how the piece is found again later. + +```bash +# UTK: an Asylum SPM run folder (summary.json or recipe.json + records/ + loops/) +python parse_utk.py ~/data/From_UTK --physical-id PDAC_COM5_01448 --out parsed + +# NLR: an XRF map of the bare film on a grid, and a DC I-V sweep over the Pt pads patterned afterwards +python parse_nlr.py ~/data/From_NLR --physical-id PDAC_COM5_01448 \ + --xrf-instrument instrument-2 --iv-instrument instrument-3 --out parsed +``` + +`parsed/` now holds a run document per run, plus any file a parser derived. Read it — it is the whole upload, +in JSON, before anything is sent. `run_document.py` states the shape. + +## Upload + +```bash +python upload_run.py parsed/*.json --account --files records +``` + +- `--files` picks which file groups go up: `records` (the per-measurement JSONs, the delivered tables, the + photographs) and `loops` (the raw arrays and plots — thousands of files, tens of minutes). `--files` with no + value uploads none and keeps the raw records in each measurement's metadata instead. +- `--dry-run` validates the documents against the ESSE schemas and stops. +- Re-running is safe: sets are found by name, samples and measurements by label, properties by what they belong + to. A second pass creates nothing. + +Then open the platform's Measurements tab: the run is there as a set, one measurement per sample. + +## Adding a lab + +Write a parser. It reads whatever that lab ships and returns the dict `run_document.py` describes — one sample +per measured position, one measurement per sample, the properties each measurement produced. Take the +measurement's workflow from standata (`standata_workflow(application, name)`); if the instrument is not in that +registry yet, add it there — `mat3ra/standata`, `assets/applications/` and `assets/workflows/` — rather than +building a workflow in Python, so the platform resolves it the same way it resolves a job's. + +Nothing in `upload_run.py` changes. + +## The notebooks + +`upload_spm_run.ipynb` and `upload_nlr_data.ipynb` run the same three steps with the same code, for people who +would rather not use a terminal. They install the pins themselves; restart the kernel after that cell if either +package was already imported. diff --git a/examples/measurement/derive_topography.py b/examples/measurement/derive_topography.py new file mode 100644 index 000000000..38430dae4 --- /dev/null +++ b/examples/measurement/derive_topography.py @@ -0,0 +1,83 @@ +#!/usr/bin/env python3 +"""AFM topography metrics already stored on the platform -> the properties they are. + +UTK's uploader kept each site's image metrics in the measurement's metadata (`frames[].imageMetrics`) +under names of its own. This reads them back and posts the ESSE properties they correspond to, so the +Results tab shows them and they are comparable with any other instrument's. Only the unambiguous +metrics are mapped; peak-to-valley, kurtosis and correlation length wait for UTK to say how they were +computed. Idempotent: a property already present on a measurement is not posted again. + + derive_topography.py [--dry-run] +""" +import argparse, os, sys, urllib.parse + +from mat3ra.api_client import APIClient + +from upload_run import account_id, base_url, find, holder + +# imageMetrics key -> the ESSE property it is. Metres in, metres out. +METRICS = { + "rq_m": {"name": "areal_surface_texture", "parameter": "Sq", "units": "m"}, + "ra_m": {"name": "areal_surface_texture", "parameter": "Sa", "units": "m"}, + "grain_radius_median_m": {"name": "grain_size", "statistic": "median", "units": "m"}, + "grain_radius_mean_m": {"name": "grain_size", "statistic": "mean", "units": "m"}, + "grain_radius_std_m": {"name": "grain_size", "statistic": "std", "units": "m"}, + "grain_radius_iqr_m": {"name": "grain_size", "statistic": "iqr", "units": "m"}, + "grain_coverage": {"name": "grain_coverage"}, +} + + +def properties_of(measurement): + """The properties one topography measurement's stored metrics stand for.""" + unit_id = measurement["workflow"]["subworkflows"][0]["units"][0]["flowchartId"] + out = [] + for frame in (measurement.get("metadata") or {}).get("frames", []): + for key, shape in METRICS.items(): + if key in frame.get("imageMetrics", {}): + out.append((unit_id, dict(shape, value=frame["imageMetrics"][key]))) + return out + + +def derive(client, set_id, dry_run): + owner_id = client.my_account.id + measurements, skip = [], 0 + while True: # the list route pages at 20 whatever the limit + page = client.measurements.list({"inSet._id": set_id, "isEntitySet": {"$ne": True}, "owner._id": owner_id}, + {"limit": 20, "skip": skip}) + measurements += page + skip += 20 + if len(page) < 20: + break + posted = present = 0 + for m in measurements: + for unit_id, prop in properties_of(m): + selector = {"source.info.origin._id": m["_id"], "data.name": prop["name"]} + for key in ("parameter", "statistic"): + if key in prop: + selector[f"data.{key}"] = prop[key] + if find(client.properties, selector, owner_id, 1): + present += 1 + continue + if not dry_run: + client.properties.create(dict(holder(prop, m["_id"], m["_sample"]["_id"], unit_id, 0), owner={"_id": owner_id})) + posted += 1 + print(f"{len(measurements)} measurements: {posted} properties {'to post' if dry_run else 'posted'}, {present} already present") + + +def main(): + ap = argparse.ArgumentParser(description=__doc__, formatter_class=argparse.RawDescriptionHelpFormatter) + ap.add_argument("set_id", help="the topography Measurement Set") + ap.add_argument("--host", default=os.environ.get("MAT3RA_HOST", "localhost:3000")) + ap.add_argument("--account", help="slug of the account the data belongs to") + ap.add_argument("--dry-run", action="store_true", help="count what would be posted and stop") + a = ap.parse_args() + url = urllib.parse.urlsplit(base_url(a.host)) + address = {"host": url.hostname, "port": url.port or (443 if url.scheme == "https" else 80), "secure": url.scheme == "https"} + client = APIClient.authenticate(**address) + if a.account: + client = APIClient.authenticate(account_id=account_id(client, a.account), **address) + derive(client, a.set_id, a.dry_run) + + +if __name__ == "__main__": + sys.exit(main()) diff --git a/examples/measurement/parse_nlr.py b/examples/measurement/parse_nlr.py new file mode 100644 index 000000000..91dcaa84c --- /dev/null +++ b/examples/measurement/parse_nlr.py @@ -0,0 +1,144 @@ +"""NLR's delivery for one piece as platform documents: the Library (the piece itself, with its layout and +synthesis), the XRF map as a Sample Set of grid points, and the DC I-V sweep as a Sample Set of pads. + +Ad hoc parser for SOF-8050: it reads the tab-separated files NLR ships and nothing else. +""" +import argparse +import json +from pathlib import Path + +from mat3ra.standata.workflows import WorkflowStandata + +from run_document import serialize + + +def standata_workflow(application_name, workflow_name): + """The procedure the instrument runs, from the standata registry — the entry the platform resolves + a job's workflow through.""" + workflow = WorkflowStandata.find_by_application_and_name(application_name, workflow_name) + if workflow is None: + raise LookupError(f"standata has no '{workflow_name}' workflow for {application_name}") + return workflow + + +def unit_id(workflow): + """The execution unit a property of this workflow comes from.""" + return workflow["subworkflows"][0]["units"][0]["flowchartId"] + +FRAME = {"origin": "wafer corner", "axes": "x, y", "units": "mm"} +DIMENSIONS = {"shape": "square", "side": 50.8, "units": "mm"} # a 2-inch substrate +IV_COLUMNS = 11 # the I-V file lists its pads row by row, eleven to a row +XRF_APPLICATION = "xrf-mapper" # the standata application whose workflow this run records + + +def read_columns(path): + """Every line of a tab-separated file after its header, split into its cells.""" + return [line.split("\t") for line in Path(path).read_text().splitlines()[1:] if line.strip()] + + +IV_APPLICATION = "probe-station" + + +def parse_nlr(folder, physical_id, xrf_instrument, iv_instrument, description="", deposition=()): + """NLR's delivery for one piece as two run documents sharing one Library: the XRF map, measured on the bare + film at the grid points, and the DC I-V sweep, measured on the Pt pads patterned afterwards. The pads are the + ones NLR probed: the I-V file lists them row by row, IV_COLUMNS to a row, so each pad gets its row and + column; where each pad sits on the wafer, and its size, come with NLR's pattern and are written onto the + same pads by label.""" + folder = Path(folder) + grid_file = sorted(folder.rglob("*xrf_grid.txt"))[0] + volts_file, amps_file = sorted(folder.rglob("IV_Volts.txt"))[0], sorted(folder.rglob("IV_Amps.txt"))[0] + run_name = grid_file.stem + images = [] # the photograph is of the piece, not of the XRF grid: it goes on the Library page + sample_set = {"name": run_name, "entitySetType": "ordered", "metadata": {}} + xrf_run_name = f"{run_name} XRF" + xrf_workflow = standata_workflow(XRF_APPLICATION, "XRF Grid Map") + xrf_unit_id = unit_id(xrf_workflow) + grid = read_columns(grid_file) + synthesis = [] + for f in deposition: # a file may hold one record or a list of them + record = json.loads(Path(f).read_text()) + synthesis.extend(record if isinstance(record, list) else [record]) + library = {"physicalId": physical_id, "name": physical_id, "description": description, + "entitySetType": "unordered", + "metadata": {"dimensions": DIMENSIONS, "frame": FRAME, "synthesis": synthesis}} + samples, xrf_measurements, xrf_properties = {}, {}, [] + for row, column, x_mm, y_mm, thickness_um, aluminium_at_pct, scandium_at_pct in grid: + label = f"r{int(row)}c{int(column)}" + # two rows for one pad would overwrite each other here and leave the row counts below + # agreeing against a dictionary that has already lost an entry + if label in samples: + raise SystemExit(f"{grid_file.name}: pad {label} appears twice") + samples[label] = {"name": f"{physical_id} {label}", "label": label, "physicalId": physical_id, + "position": {"coordinates": [float(x_mm), float(y_mm)], "units": "mm"}, + "metadata": {"row": int(row), "column": int(column)}} + xrf_measurements[label] = {"name": f"{xrf_run_name} {label}", "_sample": None, "workflow": xrf_workflow, + "setup": {"name": xrf_instrument}, "status": "finished", "_records": [], + "metadata": {"row": int(row), "column": int(column), "thickness_um": float(thickness_um), + "al_at_pct": float(aluminium_at_pct), "sc_at_pct": float(scandium_at_pct)}} + # NLR's columns become what ESSE already defines: composition is one elemental_ratio per + # element, a fraction, not a property named after the element + xrf_properties += [(label, xrf_unit_id, {"name": "film_thickness", "value": float(thickness_um), "units": "um"}, 0), + (label, xrf_unit_id, {"name": "elemental_ratio", "element": "Al", "value": float(aluminium_at_pct) / 100}, 0), + (label, xrf_unit_id, {"name": "elemental_ratio", "element": "Sc", "value": float(scandium_at_pct) / 100}, 0)] + iv_run_name = f"{physical_id} DC IV" + iv_workflow = standata_workflow(IV_APPLICATION, "DC I-V Sweep") + iv_unit_id = unit_id(iv_workflow) + volts = [[float(v) for v in cells] for cells in read_columns(volts_file)] + amps = [[float(a) for a in cells] for cells in read_columns(amps_file)] + if len(volts) != len(amps): + raise SystemExit(f"{volts_file.name}/{amps_file.name}: {len(volts)} and {len(amps)} rows") + if len(volts) % IV_COLUMNS: + raise SystemExit(f"{len(volts)} I-V rows do not fill rows of {IV_COLUMNS} pads") + for row, (bias_row, current_row) in enumerate(zip(volts, amps)): + if len(bias_row) != len(current_row): + raise SystemExit(f"row {row}: {len(bias_row)} bias points but {len(current_row)} current points") + # the pads NLR probed, by their row and column in the file; positions and sizes come with NLR's pattern + layout = [{"label": f"pad_r{i // IV_COLUMNS}c{i % IV_COLUMNS:02d}", "row": i // IV_COLUMNS, "column": i % IV_COLUMNS, + "position": None, "extent": None, "stack": None} for i in range(len(volts))] + library["metadata"]["layout"] = layout + iv_setup = {"name": iv_instrument, "settings": {"v_min": min(volts[0]), "v_max": max(volts[0]), "points": len(volts[0])}} + pads, iv_measurements, iv_properties = {}, {}, [] + for index, (pad, bias, current) in enumerate(zip(layout, volts, amps)): + label = pad["label"] + pads[label] = {"name": f"{physical_id} {label}", "label": label, "physicalId": physical_id, + "metadata": {"site": "pad", "row": pad["row"], "column": pad["column"]}} + iv_measurements[label] = {"name": f"{iv_run_name} {label}", "_sample": None, "workflow": iv_workflow, + "setup": iv_setup, "status": "finished", "_records": [], "metadata": {"row_index": index}} + iv_properties.append((label, iv_unit_id, {"name": "current_voltage_curve", "xAxis": {"label": "voltage", "units": "V"}, + "yAxis": {"label": "current", "units": "A"}, + "xDataArray": bias, "yDataSeries": [current]}, 0)) + return [{"physicalId": physical_id, "library": library, "run": xrf_run_name, "sample_set": sample_set, "images": images, + "samples": samples, "measurement_set": {"name": xrf_run_name, "entitySetType": "ordered", "metadata": {}}, + "measurements": xrf_measurements, "files": {}, "set_files": [(grid_file.name, grid_file)], + "records": grid, "properties": xrf_properties}, + {"physicalId": physical_id, "library": library, "run": iv_run_name, + "sample_set": {"name": f"{physical_id} pads", "entitySetType": "ordered", "metadata": {}}, "images": [], + "samples": pads, "measurement_set": {"name": iv_run_name, "entitySetType": "ordered", "metadata": {}}, + "measurements": iv_measurements, "files": {}, + "set_files": [(volts_file.name, volts_file), (amps_file.name, amps_file)], + "records": volts, "properties": iv_properties}] + + + +def main(): + """Read NLR's delivery and write its run document.""" + ap = argparse.ArgumentParser(description=__doc__, formatter_class=argparse.RawDescriptionHelpFormatter) + ap.add_argument("folder") + ap.add_argument("--physical-id", required=True, help="the identifier written on the physical piece, e.g. PDAC_COM5_01448") + ap.add_argument("--xrf-instrument", required=True, help="identity of the machine the grid was mapped on") + ap.add_argument("--iv-instrument", required=True, help="identity of the machine the sweep was measured on") + ap.add_argument("--description", default="", help="what the piece is, for the Library") + ap.add_argument("--deposition", nargs="*", default=[], metavar="JSON", help="deposition record(s) for the Library's synthesis") + ap.add_argument("--out", default="parsed", help="directory for the run documents (default: parsed/)") + a = ap.parse_args() + + for parsed in parse_nlr(a.folder, a.physical_id, a.xrf_instrument, a.iv_instrument, a.description, a.deposition): + print(f"{parsed['physicalId']}: {len(parsed['samples'])} samples · run {parsed['run']}: " + f"{len(parsed['measurements'])} measurements · {len(parsed['records'])} rows -> " + f"{len(parsed['set_files'])} files · {len(parsed['images'])} image(s) · {len(parsed['properties'])} properties") + print("run document:", serialize(parsed, a.out, name=parsed["run"].replace(" ", "_") + ".json")) + + +if __name__ == "__main__": + main() diff --git a/examples/measurement/parse_utk.py b/examples/measurement/parse_utk.py new file mode 100644 index 000000000..abdef8e02 --- /dev/null +++ b/examples/measurement/parse_utk.py @@ -0,0 +1,348 @@ +"""A UTK SS-PFM run folder as platform documents: a sample set of the pads, one measurement per pad, and one +hysteresis-loop property per pad — the eight loops combined, with the loop parameters (mean, population standard +deviation, count) inside it. Individual loops stay in the measurement's files. + +Ad hoc parser for SOF-8050: it reads the shape UTK's afm-lib writes and nothing else. +""" +import argparse, ast, json, math, re, statistics, struct +from datetime import datetime, timezone +from pathlib import Path + +from mat3ra.standata.workflows import WorkflowStandata + +from run_document import serialize + + +def standata_workflow(application_name, workflow_name): + """The procedure the instrument runs, from the standata registry — the entry the platform resolves + a job's workflow through.""" + workflow = WorkflowStandata.find_by_application_and_name(application_name, workflow_name) + if workflow is None: + raise LookupError(f"standata has no '{workflow_name}' workflow for {application_name}") + return workflow + + +def unit_id(workflow): + """The execution unit a property of this workflow comes from.""" + return workflow["subworkflows"][0]["units"][0]["flowchartId"] + + + +FIELD = {"off_field": "off", "on_field": "on"} +# loop_params key -> parameters path (units: voltages in xAxis.units, responses in yAxis.units) +PARAMETERS = { + "imprint_v": ("imprint",), + "v_c_rising": ("coerciveVoltage", "rising"), + "v_c_falling": ("coerciveVoltage", "falling"), + "loop_width_v": ("loopWidth",), + "loop_height_m": ("loopHeight",), + "remnant_rising_m": ("remanentResponse", "rising"), + "remnant_falling_m": ("remanentResponse", "falling"), +} +INSTRUMENT_NAME = "asylum-spm" # the standata application whose workflow this run records +g = lambda v: float(f"{v:.6g}") + + +def load_npy(path): + """A one-dimensional float32/float64 .npy file as a list of floats (no numpy: the format is a header + raw values).""" + data = Path(path).read_bytes() + if data[:6] != b"\x93NUMPY": + raise ValueError(f"{path}: not a .npy file") + header_length = struct.unpack(" 1 else 0.0, "count": len(values)} + + +def combine_pad(label, records, run_dir): + """One hysteresis_loop property for a pad: point-wise mean of the response over its loops (on and off), and + each loop parameter as mean / spread / count over the loops. None when no loop has curves.""" + bias, series, params = None, {"on": [], "off": []}, {"on": {}, "off": {}} + for r in records: + loops_dir = run_dir / "loops" / (r.get("out_stem") or Path(r["file_path"]).stem) + for branch, field in FIELD.items(): + lp = r["loop_params"].get(branch) or {} + for key, path in PARAMETERS.items(): + if lp.get(key) is not None: + params[field].setdefault(path, []).append(lp[key]) + bias_p = loops_dir / f"bias_{field}.npy" + if not bias_p.exists() or lp.get("phase_offset_deg") is None: + continue + b = load_npy(bias_p) + bias = bias or b + curve = response_curve(loops_dir, field, lp["phase_offset_deg"], b) + # averaged point by point, so a loop measured on a different bias axis is dropped rather + # than folded in: same length is not the same voltages + if curve is not None and same_axis(b, bias): + series[field].append(curve) + if bias is None or not series["on"] or not series["off"]: + return None + parameters = {} + for field in ("on", "off"): + block = {} + for path, vals in params[field].items(): + node = block + for k in path[:-1]: + node = node.setdefault(k, {}) + node[path[-1]] = parameter_statistics(vals) + parameters[field] = block + return {"name": "hysteresis_loop", "legend": ["on", "off"], + "xAxis": {"label": "bias", "units": "V"}, "yAxis": {"label": "response", "units": "m"}, + "xDataArray": [g(v) for v in bias], "yDataSeries": [pointwise_mean(series["on"]), pointwise_mean(series["off"])], + "parameters": parameters} + +RECORD_GROUPS = ("labels", "file_path", "requested_params", "instrument_params", "loop_params", "channel_stats") +SAMPLE_KEY = "labels.site_label" # a record's reference to its sample — stays on every record + + +def _flatten(d, prefix=""): + """Nested dict → one level, keys joined with dots.""" + out = {} + for k, v in (d or {}).items(): + key = f"{prefix}.{k}" if prefix else k + if isinstance(v, dict): + out.update(_flatten(v, key)) + else: + out[key] = v + return out + + +def _unflatten(flat): + """The inverse of _flatten.""" + out = {} + for key, v in flat.items(): + node = out + parts = key.split(".") + for part in parts[:-1]: + node = node.setdefault(part, {}) + node[parts[-1]] = v + return out + + +def _constant_keys(flat_rows): + """Keys present in every row with the same value everywhere. A key one row lacks is not constant, even if the + rows that have it agree — otherwise a group that is null on one record and absent on another looks shared.""" + if not flat_rows: + return set() + keys = set.intersection(*(set(row) for row in flat_rows)) + return {k for k in keys if all(row[k] == flat_rows[0][k] for row in flat_rows)} + + +def factor_records(records): + """Each record keeps only what is unique to it. Fields identical across the whole run move to the measurement + (`common`); fields identical across one sample's records move to that sample's metadata; the record keeps its + sample reference and the values that actually vary per loop (measured 128 → 43 / 4 / 81 on the 704-record run).""" + if not records: + return {}, {}, [] + flat = [_flatten({k: r.get(k) for k in RECORD_GROUPS}) for r in records] + common_keys = _constant_keys(flat) - {SAMPLE_KEY} + by_sample = {} + for row in flat: + by_sample.setdefault(row[SAMPLE_KEY], []).append(row) + sample_keys = set.intersection(*(_constant_keys(rows) for rows in by_sample.values())) - common_keys - {SAMPLE_KEY} + # a group can be a dict on one record and null on another: read with .get, never index + common = _unflatten({k: flat[0].get(k) for k in sorted(common_keys)}) + per_sample = {label: _unflatten({k: rows[0].get(k) for k in sorted(sample_keys)}) for label, rows in by_sample.items()} + slim = [_unflatten({k: v for k, v in row.items() if k not in common_keys and k not in sample_keys}) for row in flat] + return common, per_sample, slim + +def starting_site(recipe): + """The site the run started from: r0c00 when the recipe has it, else the first one listed.""" + return next((s for s in recipe["sites"] if s["label"] == "r0c00"), recipe["sites"][0]) + + +def thinned_curves(prop, points=12): + """The curves at `points` evenly spaced samples: an ESSE example shows the shape, not the data.""" + x = prop["xDataArray"] + if len(x) <= points: + return {} + keep = [round(i * (len(x) - 1) / (points - 1)) for i in range(points)] + return {"xDataArray": [x[i] for i in keep], + "yDataSeries": [[s[i] for i in keep] for s in prop["yDataSeries"]]} + + +def registration(recipe): + """The instrument's frame as UTK stated it: one anchor in words (from recipe.context) and where the run started + on the stage. Recorded, not interpreted.""" + ctx = recipe.get("context", "") + m = re.search(r"starting point is (.+?)(?:\.|$)", ctx) + r0 = starting_site(recipe) + return {"frame": "asylum-spm stage", "units": "m", "anchor": m.group(1).strip() if m else ctx, + f"{r0['label']}_stage_m": [r0["x_stage_m"], r0["y_stage_m"]]} + + +def sample_files(label, records, run_dir, slim_by_index): + """Files of one sample's measurement: one JSON per record (the fields unique to it) and, when the run folder has them, + the loop arrays and annotated plots. Returned as (relative name, payload) where payload is text or a Path.""" + out = [] + for r, slim in zip(records, slim_by_index): + step, point = r["labels"].get("step_index", 0), r["labels"].get("point_index", 0) + out.append((f"records/step{step}_pt{point:02d}.json", json.dumps(slim, indent=1))) + d = run_dir / "loops" / r.get("out_stem", "") + if r.get("out_stem") and d.is_dir(): + for f in sorted(d.iterdir()): + if f.suffix in (".npy", ".png"): + out.append((f"loops/{d.name}/{f.name}", f)) + return out + + +def parse(run_dir, physical_id, limit_records=None, deposition=None, instrument="instrument-1"): + """The whole run folder as platform documents: sample set, samples, measurement set, one measurement per sample, files, one loop property per fully measured sample.""" + run_dir = Path(run_dir) + recipe, session, all_records = load_run(run_dir) + records = all_records[:limit_records] if limit_records else all_records + if not physical_id: + raise ValueError("--physical-id: the identifier written on the physical piece is required") + run_name = session.get("name") or run_dir.name + reg = registration(recipe) + sample_set = {"name": run_name, "entitySetType": "ordered", "metadata": {}} + # the Library: the piece itself. NLR's deposition record(s), when UTK dropped them into the run folder as + # deposition*.json or --deposition names them, are its synthesis. + deposition_files = [Path(deposition)] if deposition else sorted(run_dir.glob("deposition*.json")) + synthesis = [] + for f in deposition_files: + d = json.loads(f.read_text()) + synthesis.extend(d if isinstance(d, list) else [d]) + library = {"physicalId": physical_id, "name": physical_id, "description": "", "entitySetType": "unordered", + "metadata": {"synthesis": synthesis}} + # the photograph of the piece: any image at the run-folder root + images = [(f.name, f) for f in sorted(run_dir.iterdir()) if f.suffix.lower() in (".jpg", ".jpeg", ".png")] + # samples in recipe order (the set is ordered; the server assigns inSet.index as they are moved in) + samples = {s["label"]: {"name": f"{physical_id} {s['label']}", "label": s["label"], "physicalId": physical_id, + "position": {"coordinates": [s["x_stage_m"], s["y_stage_m"]], "units": "m"}, + "metadata": {"registration": reg}} + for s in recipe["sites"]} + if limit_records: + # a trial run must be a prefix of a full one: only samples whose records ALL made the cut, so no sample is + # ever published with a partial loop count that a full run would then skip as "already there" + full_counts, kept_counts = {}, {} + for r in all_records: + full_counts[r["labels"]["site_label"]] = full_counts.get(r["labels"]["site_label"], 0) + 1 + for r in records: + kept_counts[r["labels"]["site_label"]] = kept_counts.get(r["labels"]["site_label"], 0) + 1 + complete = {label for label, n in kept_counts.items() if n == full_counts[label]} + records = [r for r in records if r["labels"]["site_label"] in complete] + samples = {label: sample for label, sample in samples.items() if label in complete} + common, per_sample, slim_records = factor_records(records) + for label, const in per_sample.items(): + if label in samples: + samples[label]["metadata"].update(const) + workflow = standata_workflow(INSTRUMENT_NAME, "SS-PFM Hysteresis Loop") + unit = unit_id(workflow) + measurement_set = {"name": run_name, "entitySetType": "ordered", + "metadata": {"session": session, "recipe": recipe, "context": recipe.get("context", ""), + "loop_settings": recipe["per_site"][0]["loop_settings"], + "sites": list(samples), "common": common, "registration": reg}} + # the setup block, Measurement : setup :: Job : compute — the machine and the sitting; the technique + # (asylum-spm, SS-PFM) is the workflow's application. The run folder does not name the machine: --instrument does. + started = session.get("started_ts") + setup_block = {"name": instrument, + "session": {k: v for k, v in {"name": session.get("name"), + "started": datetime.fromtimestamp(started, timezone.utc).isoformat().replace("+00:00", "Z") if started else None}.items() if v}, + **({"settings": common["instrument_params"]} if common.get("instrument_params") else {})} + by_sample, slim_by_sample = {}, {} + for r, slim in zip(records, slim_records): + lab = r["labels"]["site_label"] + by_sample.setdefault(lab, []).append(r); slim_by_sample.setdefault(lab, []).append(slim) + # one measurement per sample: Measurement : Sample :: Job : Material + measurements, files, properties, skipped = {}, {}, [], [] + for label in samples: + recs = by_sample.get(label, []) + measurements[label] = {"name": f"{run_name} {label}", "_sample": None, "workflow": workflow, + "setup": setup_block, "status": "finished", + "metadata": {"run_dir": session.get("run_dir", str(run_dir)), "recordsCount": len(recs)}} + slim_by_label = slim_by_sample.get(label, []) + measurements[label]["_records"] = slim_by_label # not sent; upload() decides files vs metadata + files[label] = sample_files(label, recs, run_dir, slim_by_sample.get(label, [])) + prop = combine_pad(label, recs, run_dir) if recs else None + (properties.append((label, unit, prop, 0)) if prop else skipped.append(label)) + return {"physicalId": physical_id, "library": library, "run": run_name, "sample_set": sample_set, "images": images, "samples": samples, + "measurement_set": measurement_set, "measurements": measurements, "files": files, "set_files": [], + "records": records, "properties": properties, "skipped": skipped} + + + +def same_axis(candidate, reference, tolerance=1e-9): + """Whether two bias axes are the same sweep: equal length and equal voltages within tolerance.""" + return len(candidate) == len(reference) and all( + abs(a - b) <= tolerance + 1e-6 * abs(b) for a, b in zip(candidate, reference)) + + +def main(): + """Read a UTK run folder and write its run document.""" + ap = argparse.ArgumentParser(description=__doc__, formatter_class=argparse.RawDescriptionHelpFormatter) + ap.add_argument("run_dir") + ap.add_argument("--physical-id", required=True, help="the identifier written on the physical piece, e.g. PDAC_COM5_01448") + ap.add_argument("--out", default="parsed", help="directory for the run document and the records cut from the run (default: parsed/)") + ap.add_argument("--limit-records", type=int, help="trial: only the first N records and the samples they belong to") + ap.add_argument("--deposition", help="NLR HTEM record (json) kept in the run's sample set metadata") + ap.add_argument("--instrument", default="instrument-1", help="identity of the machine the run was measured on (the run folder does not record it)") + ap.add_argument("--emit-example", help="write the property with the most loops to this path — the ESSE example") + a = ap.parse_args() + + parsed = parse(a.run_dir, a.physical_id, a.limit_records, a.deposition, a.instrument) + files = sum(len(v) for v in parsed["files"].values()) + print(f"{parsed['physicalId']}: {len(parsed['samples'])} samples · run {parsed['run']}: " + f"{len(parsed['measurements'])} measurements · {len(parsed['records'])} records -> {files} files · " + f"{len(parsed['properties'])} samples with a combined loop" + + (f" · no curves: {len(parsed['skipped'])} samples" if parsed["skipped"] else "")) + if a.emit_example and parsed["properties"]: + label, _, prop, _rep = max(parsed["properties"], key=lambda t: t[2]["parameters"]["off"].get("imprint", {}).get("count", 0)) + Path(a.emit_example).write_text(json.dumps(dict(prop, **thinned_curves(prop)), indent=4) + "\n") + print(f"example written from sample {label} -> {a.emit_example}") + print("run document:", serialize(parsed, a.out)) + + +if __name__ == "__main__": + main() diff --git a/examples/measurement/requirements.txt b/examples/measurement/requirements.txt new file mode 100644 index 000000000..8d3694456 --- /dev/null +++ b/examples/measurement/requirements.txt @@ -0,0 +1,11 @@ +# The measurement upload: pip install -r requirements.txt +# +# api-client, standata and esse are pinned to the branch carrying the REST endpoints, the instrument +# registry entries and the Sample/Measurement schemas. They become version pins when it releases. + +mat3ra-api-client @ git+https://github.com/mat3ra/api-client.git@feature/SOF-8051 +mat3ra-standata @ git+https://github.com/mat3ra/standata.git@feature/SOF-8051 +mat3ra-esse @ git+https://github.com/mat3ra/esse.git@feature/SOF-8051 +mat3ra-notebooks-utils[api] +ipywidgets +requests diff --git a/examples/measurement/run_document.py b/examples/measurement/run_document.py new file mode 100644 index 000000000..30ca4d032 --- /dev/null +++ b/examples/measurement/run_document.py @@ -0,0 +1,73 @@ +"""The run document: what a parser writes and the uploader takes. + +One JSON file per run, beside the files it names. Everything instrument-specific has already happened by the +time it exists — reading the lab's delivery is `parse_utk.py` / `parse_nlr.py`, uploading it is `upload_run.py`, +and neither knows anything about the other. + + { + "physicalId": "PDAC_COM5_01448", the identifier written on the physical piece + "run": "com5_1448_xrf_grid XRF", names the measurement set + "sample_set": {name, entitySetType, metadata}, + "samples": {label: sample document}, + "measurement_set": {name, entitySetType, metadata}, + "measurements": {label: measurement document}, _sample is filled in at upload + "properties": [[label, unitId, property document, repetition]], + "set_files": [[name, path]], the set's own files: a grid, a photograph of the piece + "files": {label: [[name, path]]}, one measurement's files; the name's first segment is its group + "skipped": [label] informational: samples the run produced no property for + } + +Paths are relative to the document, so a run folder moves as a whole. A parser that derives a file (a record +JSON it cut from a larger one) writes it here too — `serialize` takes text in place of a path and stores it. +""" +import json +from pathlib import Path + + +def serialize(parsed, out_dir, name="run.json"): + """The parsed run as a document on disk. Text payloads are written out; Paths are recorded relative to it.""" + out_dir = Path(out_dir) + out_dir.mkdir(parents=True, exist_ok=True) + + def store(prefix, entries): + out = [] + for file_name, payload in entries: + if isinstance(payload, Path): + out.append([file_name, relative(payload, out_dir)]) + else: # derived text: this document owns it + target = out_dir / "files" / prefix / file_name + target.parent.mkdir(parents=True, exist_ok=True) + target.write_text(payload) + out.append([file_name, relative(target, out_dir)]) + return out + + document = {k: v for k, v in parsed.items() if k not in ("files", "set_files", "images", "measurements")} + document["measurements"] = {label: {k: v for k, v in m.items() if k != "_records"} + for label, m in parsed["measurements"].items()} + document["records_by_sample"] = {label: m["_records"] for label, m in parsed["measurements"].items() if m.get("_records")} + document["set_files"] = store("set", list(parsed.get("set_files", [])) + list(parsed.get("images", []))) + document["files"] = {label: store(label, entries) for label, entries in parsed.get("files", {}).items()} + document["properties"] = [list(p) for p in parsed.get("properties", [])] + path = out_dir / name + path.write_text(json.dumps(document, indent=1) + "\n") + return path + + +def relative(path, out_dir): + """`path` as the document will name it: relative when it can be, absolute when it lives elsewhere.""" + path, out_dir = Path(path).resolve(), Path(out_dir).resolve() + try: + return path.relative_to(out_dir).as_posix() + except ValueError: + return path.as_posix() + + +def load(path): + """A run document with its file paths resolved against its own location.""" + path = Path(path) + document = json.loads(path.read_text()) + resolve = lambda entries: [(name, (path.parent / file_path).resolve()) for name, file_path in entries] + document["set_files"] = resolve(document.get("set_files", [])) + document["files"] = {label: resolve(entries) for label, entries in document.get("files", {}).items()} + document["properties"] = [tuple(p) for p in document.get("properties", [])] + return document diff --git a/examples/measurement/upload_nlr_data.ipynb b/examples/measurement/upload_nlr_data.ipynb new file mode 100644 index 000000000..bf86f6df7 --- /dev/null +++ b/examples/measurement/upload_nlr_data.ipynb @@ -0,0 +1,244 @@ +{ + "cells": [ + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "# Overview\n", + "\n", + "This example uploads the data NLR delivers for one physical piece: the folder becomes a Sample Set with one Sample per measured pad, and two runs over those same pads \u2014 an XRF map and a DC I-V sweep \u2014 each a Measurement Set with one Measurement per Sample and its Setup, the delivered tables as files, and one Property per measured quantity.\n", + "A delivered folder holds the XRF grid table, the DC I-V tables and photographs of the piece, and re-running the notebook adds only what is missing." + ] + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "## Install the API client\n", + "\n", + "The samples, measurements and files endpoints are not released yet, so the client is installed from its branch until it merges. Restart the kernel after this cell." + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "%pip install -r requirements.txt" + ] + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "## Set Parameters\n", + "\n", + "- **HOST**: platform the data is uploaded to\n", + "- **DATA_DIR**: the delivered folder beside this notebook \u2014 the XRF grid table, the DC I-V tables and the photographs\n", + "- **PHYSICAL_ID**: the identifier written on the physical piece the measured pads are part of \u2014 every Sample carries it\n", + "- **XRF_INSTRUMENT**: the machine the XRF map was measured on\n", + "- **IV_INSTRUMENT**: the machine the DC I-V sweep was measured on\n", + "- **ACCOUNT_SLUG**: account the data belongs to, empty for the default account\n", + "- **FILES**: which files to upload" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "import urllib.parse\n", + "\n", + "HOST = \"https://alphafilm.mat3ra.com\"\n", + "DATA_DIR = \"nlr\"\n", + "PHYSICAL_ID = \"\"\n", + "XRF_INSTRUMENT = \"\"\n", + "IV_INSTRUMENT = \"\"\n", + "ACCOUNT_SLUG = \"\"\n", + "FILES = [\"records\"] # the delivered tables and photographs; [] uploads none\n", + "\n", + "url = urllib.parse.urlsplit(HOST)\n", + "address = {\n", + " \"host\": url.hostname,\n", + " \"port\": url.port or (443 if url.scheme == \"https\" else 80),\n", + " \"secure\": url.scheme == \"https\",\n", + "}" + ] + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "## Authenticate and initialize API client\n", + "\n", + "### Authenticate\n", + "Authenticate in the browser (OIDC device flow) or via JupyterLite host injection. Credentials are stored in environment variables.\n", + "\n", + "### Initialize API client\n", + "Create an authenticated API client and resolve the owner account ID." + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "from mat3ra.notebooks_utils.auth import authenticate\n", + "\n", + "await authenticate()" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "from mat3ra.api_client import APIClient\n", + "\n", + "client = APIClient.authenticate(**address)" + ] + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "# Imports" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "from pathlib import Path\n", + "\n", + "from parse_nlr import parse_nlr\n", + "from run_document import load, serialize\n", + "from upload_run import account_id, upload" + ] + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "## Parse the delivered folder\n", + "\n", + "Read the folder into the documents the platform stores: one Sample Set, and one run per technique over its Samples. Nothing is uploaded yet." + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "# reading the delivery and writing one run document per technique: nothing here talks to the platform\n", + "runs = []\n", + "for parsed in parse_nlr(Path(DATA_DIR), PHYSICAL_ID, XRF_INSTRUMENT, IV_INSTRUMENT):\n", + " document_path = serialize(parsed, \"parsed\", name=parsed[\"run\"].replace(\" \", \"_\") + \".json\")\n", + " run = load(document_path)\n", + " runs.append(run)\n", + " print(\n", + " f\"{run['physicalId']}: {len(run['samples'])} samples (ordered set) \u00b7 run {run['run']}: \"\n", + " f\"{len(run['measurements'])} measurements (ordered set, one per sample) \u00b7 \"\n", + " f\"{len(run['set_files'])} files \u00b7 {len(run['properties'])} properties \u00b7 {document_path}\"\n", + " )" + ] + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "## Select the account\n", + "\n", + "`ACCOUNT_SLUG` re-authenticates the client against that account, so the run is read and written there." + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "if ACCOUNT_SLUG:\n", + " client = APIClient.authenticate(account_id=account_id(client, ACCOUNT_SLUG), **address)" + ] + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "## Upload the data\n", + "\n", + "Create the Sample Set and its Samples, then a Measurement Set per technique with one Measurement per Sample, the delivered files and the Properties." + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "for run in runs:\n", + " upload(client, run, files=FILES)" + ] + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "## Find the data in the web app\n", + "\n", + "Each technique is a folder in the account's Measurements tab." + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "print(f\"Open {HOST}, your account's Measurements tab: {', '.join(run['run'] for run in runs)}\")" + ] + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "## References\n", + "\n", + "- [Mat3ra REST API](https://docs.mat3ra.com/rest-api/overview/)" + ] + } + ], + "metadata": { + "colab": { + "name": "upload_nlr_data.ipynb", + "provenance": [] + }, + "kernelspec": { + "display_name": "Python 3", + "language": "python", + "name": "python3" + }, + "language_info": { + "codemirror_mode": { + "name": "ipython", + "version": 3 + }, + "file_extension": ".py", + "mimetype": "text/x-python", + "name": "python", + "nbconvert_exporter": "python", + "pygments_lexer": "ipython3", + "version": "3.8.6" + } + }, + "nbformat": 4, + "nbformat_minor": 1 +} diff --git a/examples/measurement/upload_run.py b/examples/measurement/upload_run.py new file mode 100755 index 000000000..2e80b7743 --- /dev/null +++ b/examples/measurement/upload_run.py @@ -0,0 +1,321 @@ +#!/usr/bin/env python3 +"""Run documents -> the platform: the sample set, its samples, the measurement set, one measurement per sample, +the files and the properties, validated against the ESSE schemas before anything is sent. + + upload_run.py parsed/run.json --account # upload it + upload_run.py parsed/*.json --account --files records loops + upload_run.py parsed/run.json --dry-run # validate only + +It knows nothing about any instrument. Reading a lab's delivery into a run document is a parser's job — +`parse_utk.py` for a UTK SS-PFM run, `parse_nlr.py` for NLR's delivery — and `run_document.py` states the +shape they agree on. A new lab is a new parser; nothing here changes. + +Requires Python 3.9+ and `pip install -r requirements.txt`. Credentials come from the environment: +OIDC_ACCESS_TOKEN, or ACCOUNT_ID + AUTH_TOKEN (an API token from Preferences); MAT3RA_HOST picks the host. +""" +import argparse, concurrent.futures, os, sys, threading, time, urllib.parse +from pathlib import Path + +import requests +from mat3ra.api_client import APIClient + +from run_document import load + +from mat3ra.esse import ESSE + +def holder(prop, measurement_id, sample_id, unit_id, repetition): + """The property holder the platform stores: the data, where it came from (measurement, sample, workflow unit) and a + repetition index — 0, since a measurement holds one sample and one loop property.""" + return {"data": prop, + "source": {"type": "external", + "info": {"origin": {"_id": measurement_id, "cls": "Measurement"}, + "subject": {"_id": sample_id, "cls": "Sample"}, + "unitId": unit_id}}, + "exabyteId": [], "repetition": repetition} + + +def validate(parsed): + """Validate every document against the ESSE schemas; returns the number of invalid ones.""" + esse = ESSE() + schemas = {x["$id"]: x for x in esse.schemas} + errors = 0 + for smp in parsed["samples"].values(): + try: + esse.validate(smp, schemas["sample"]) + except Exception as e: + errors += 1; print("SAMPLE INVALID", smp["label"], str(e)[:200]) + for label, m in parsed["measurements"].items(): + m = {k: v for k, v in m.items() if k != "_records"}; m["_sample"] = {"_id": "dryrun", "cls": "Sample", "slug": label} + try: + esse.validate(m, schemas["measurement"]) + except Exception as e: + errors += 1; print("MEASUREMENT INVALID", label, str(e)[:300]); break + unvalidated = set() + for label, uid, prop, rep in parsed["properties"]: + # by the property's own name: NLR's thickness, atomic fractions and I-V curve are not the + # hysteresis loop, and ESSE has no schema for them yet, so they are reported, not failed + # by the property's own name, wherever the directory keeps it: scalar, non-scalar, structural + suffix = "/" + prop["name"].replace("_", "-") + schema = next((s for schema_id, s in schemas.items() + if schema_id.startswith("properties-directory/") and schema_id.endswith(suffix)), None) + if schema is None: + unvalidated.add(prop["name"]); continue + try: + esse.validate(prop, schema) + esse.validate(holder(prop, "dryrun", "dryrun", uid, rep), schemas["property/holder"]) + except Exception as e: + errors += 1; print("PROPERTY INVALID", label, str(e)[:300]) + if unvalidated: + print("no ESSE schema yet, not validated:", ", ".join(sorted(unvalidated))) + return errors + + + +def base_url(host): + """`https://` unless told otherwise: a bare hostname becomes https, localhost/127.0.0.1 http, a URL is kept.""" + host = host.rstrip("/") + if host.startswith(("http://", "https://")): + return host + return ("http://" if host.split(":")[0] in ("localhost", "127.0.0.1") else "https://") + host + + +def find(endpoint, query, owner_id, limit=100): + """The account's entities matching `query` — a query bypasses the route's account scoping, so every lookup is + narrowed. A page as long as the limit is refused: a missed member means a duplicate on the next run.""" + found = endpoint.list(dict(query, **{"owner._id": owner_id}), {"limit": limit}) + if len(found) >= limit and limit > 1: + raise SystemExit(f"{endpoint.name}: more than {limit} match — this script cannot page yet; stop rather than duplicate") + return found + + +def account_id(client, slug_or_name): + """`--account ` → the id of that account: the slug the platform shows it under, else its display name.""" + accounts = client.list_accounts() + for field in ("slug", "name"): + account = next((a for a in accounts if a.get(field) == slug_or_name), None) + if account: + return account["_id"] + raise SystemExit(f"account '{slug_or_name}' is not one of yours: " + f"{', '.join(a.get('slug') or a['name'] for a in accounts)}") + + +THREAD = threading.local() + + +def thread_client(client): + """A client of the calling thread's own: the endpoints keep the last response on the connection they share, so + the file workers cannot use one between them. Building one costs nothing — no request is made until it is used.""" + if not hasattr(THREAD, "client"): + THREAD.client = APIClient(host=client.host, port=client.port, version=client.version, secure=client.secure, + auth=client.auth, timeout_seconds=client.timeout_seconds) + return THREAD.client + + +def put_file(client, name, payload, owner_id): + """One file into the account's file store: a file on disk through a signed PUT, text through the files route. + Eight connections at once is enough to make a name resolution fail now and then, so a refused connection is + tried again twice before it takes the run down with it.""" + for attempt in range(3): + try: + if isinstance(payload, Path): + return client.files.put(payload, name, account_id=owner_id) + return client.files.create(name, payload, account_id=owner_id) + except (requests.exceptions.ConnectionError, requests.exceptions.Timeout): + if attempt == 2: + raise + time.sleep(2 * (attempt + 1)) + + +REPLACED_WHOLE = ("layout", "dimensions", "frame") # a layout is one thing, not a list that grows + + +def merge_metadata(existing, incoming): + """`incoming` on top of `existing`, keeping what neither replaces. A list grows by the entries it + does not already hold - a re-upload brings deposition records the set has never seen, + and taking only absent keys would drop them because `deposition` is already there.""" + merged = dict(existing) + for key, value in incoming.items(): + held = merged.get(key) + if key in REPLACED_WHOLE: + merged[key] = value + elif isinstance(held, list) and isinstance(value, list): + merged[key] = held + [v for v in value if v not in held] + elif isinstance(held, dict) and isinstance(value, dict): + merged[key] = merge_metadata(held, value) + else: + merged[key] = value + return merged + + +def ensure_set(endpoint, doc, owner_id, parent_id=None): + """The set with this name in the account, created when missing; returns (set, created). + An existing set takes any metadata it does not have yet - a later upload may carry a + deposition record the set was created without - and is moved under `parent_id` when it + is not there already.""" + candidates = find(endpoint, {"isEntitySet": True, "name": doc["name"]}, owner_id, 20) + if parent_id: + # the set is this run's only if it sits under this Library, or under no Library at all (an + # upload from before Libraries existed, adopted now); one under another piece is another run + libraries = {s["_id"] for s in find(endpoint, {"isEntitySet": True, "physicalId": {"$exists": True}}, owner_id, 100)} + def parents(candidate): + return {ref.get("_id") for ref in candidate.get("inSet", [])} + candidates = [c for c in candidates if parent_id in parents(c) or not (parents(c) & libraries)] + if not candidates: + body = dict(doc, owner={"_id": owner_id}) + if parent_id: + body["parentSetId"] = parent_id + return endpoint.create_set(body), True + + existing = candidates[0] + changes = {} + merged = merge_metadata(existing.get("metadata") or {}, doc.get("metadata") or {}) + if merged != (existing.get("metadata") or {}): + changes["metadata"] = merged + for key in ("description", "physicalId"): # a Library's own fields, filled in when the set predates them + if doc.get(key) and not existing.get(key): + changes[key] = doc[key] + if changes: + endpoint.update_set(existing["_id"], changes) + existing = dict(existing, **changes) + if parent_id and parent_id not in {s.get("_id") for s in existing.get("inSet", [])}: + endpoint.move_to_set(existing["_id"], None, parent_id) + return existing, False + + +FILE_GROUPS = ("records", "loops") + + +def run_files(parsed, groups=("records",)): + """Every file this run puts, as (name, payload) pairs. `groups` names what to include per + measurement: "records" for the record JSONs, "loops" for the loop arrays and plots, which are + large. The set's own files - NLR's grid, the photographs of the piece - are few and always go. + `name` carries the label it belongs to so `destination` can address it.""" + unknown = set(groups) - set(FILE_GROUPS) + if unknown: + raise SystemExit(f"unknown file group(s): {', '.join(sorted(unknown))}; choose from {', '.join(FILE_GROUPS)}") + files = [(f"set/{name}", payload) for name, payload in parsed["set_files"]] + for label, file_list in parsed["files"].items(): + files += [(f"{label}/{name}", payload) for name, payload in file_list + if ("loops" if name.startswith("loops/") else "records") in groups] + return files + + +def destination(name, set_id, measurement_ids): + """Where one of `run_files`' entries is put: the set's own, or the measurement it belongs to.""" + label, _, rest = name.partition("/") + if label == "set": + return f"sets/{set_id}/{rest}" + return f"measurements/{measurement_ids[label]}/{rest}" + + +def upload(client, parsed, files=("records",), properties=True): + """The run onto its Sample Set: the set, its samples, the measurement set, one measurement per + sample, the files and the properties. `files` names which groups to upload - see FILE_GROUPS; + an empty list uploads none and keeps the raw records in each measurement's metadata instead. + Idempotent: sets by run name, members by name/label; files re-put; properties posted only when + missing. `properties=False` uploads the run without them - the platform rejects a property + whose name ESSE has no schema for, and the files still carry the data to derive them from.""" + groups = list(files or []) + if not properties and groups and "loops" not in groups and any( + name.startswith("loops/") for file_list in parsed.get("files", {}).values() for name, _ in file_list): + # the properties are derived from the loop arrays; with no properties posted the arrays are the + # only record of them, so they go up whatever the groups asked for + groups.append("loops") + print("properties skipped: uploading the loop arrays too, so they can be derived later") + uploads = run_files(parsed, groups) if groups else [] + run_name = parsed["run"] + owner = {"_id": client.my_account.id} + # the Library: the physical piece, a set named by its physicalId that every Sample Set measured + # on it belongs to; its metadata carries the layout and synthesis + library_id = None + if parsed.get("library"): + library, created_library = ensure_set(client.samples, parsed["library"], owner["_id"]) + library_id = library["_id"] + print(f"library {library_id} ({library['name']}{', created' if created_library else ''})") + sample_set, created_set = ensure_set(client.samples, parsed["sample_set"], owner["_id"], library_id) + set_id = sample_set["_id"] + in_set = {s.get("label"): s for s in find(client.samples, {"inSet._id": set_id, "isEntitySet": {"$ne": True}}, owner["_id"], 500)} + sample_ids, created = {}, 0 + for label, sample_doc in parsed["samples"].items(): + if label in in_set: + sample_ids[label] = in_set[label]["_id"] + continue + doc = client.samples.create(dict(sample_doc, owner=owner)) + client.samples.move_to_set(doc["_id"], None, set_id) + sample_ids[label] = doc["_id"] + created += 1 + print(f"sample set {set_id} ({sample_set['name']}, ordered{', created' if created_set else ''}): {len(sample_ids)} samples, {created} created") + measurement_set, created_measurement_set = ensure_set(client.measurements, parsed["measurement_set"], owner["_id"]) + existing = {m["name"]: m for m in find(client.measurements, {"inSet._id": measurement_set["_id"], "isEntitySet": {"$ne": True}}, owner["_id"], 500)} + measurement_ids, measurements_created = {}, 0 + for label, measurement_doc in parsed["measurements"].items(): + if measurement_doc["name"] in existing: + measurement_ids[label] = existing[measurement_doc["name"]]["_id"] + continue + body = dict(measurement_doc) + body["_sample"] = {"_id": sample_ids[label], "cls": "Sample"} + records = parsed.get("records_by_sample", {}).get(label) + if not uploads and records: # no file store: keep the raw records in the measurement's metadata + body["metadata"] = dict(body["metadata"], records=records) + doc = client.measurements.create(dict(body, owner=owner)) + client.measurements.move_to_set(doc["_id"], None, measurement_set["_id"]) + measurement_ids[label] = doc["_id"] + measurements_created += 1 + print(f"measurement set {measurement_set['_id']} ({run_name}, ordered{', created' if created_measurement_set else ''}): " + f"{len(measurement_ids)} measurements, {measurements_created} created") + if uploads: + # one request per file (~1-2 s each), eight in flight + jobs = [(destination(name, set_id, measurement_ids), payload) for name, payload in uploads] + with concurrent.futures.ThreadPoolExecutor(max_workers=8) as pool: + for done, _ in enumerate(pool.map(lambda job: put_file(thread_client(client), *job, owner["_id"]), jobs), 1): + if done % 200 == 0: + print(f" files: {done}/{len(jobs)}", flush=True) + print(f"files: {len(jobs)} put") + if not properties: + print(f"properties: skipped ({len(parsed['properties'])} not posted)") + return + posted = 0 + for label, unit_id, prop, repetition in parsed["properties"]: # properties/create is not idempotent: skip what is there + selector = {"source.info.origin._id": measurement_ids[label], "data.name": prop["name"], "repetition": repetition} + if "element" in prop: # one measurement holds an elemental_ratio per element + selector["data.element"] = prop["element"] + present = find(client.properties, selector, owner["_id"], 1) + if present: + continue + client.properties.create(dict(holder(prop, measurement_ids[label], sample_ids[label], unit_id, repetition), owner=owner)) + posted += 1 + print(f"properties: {posted} posted, {len(parsed['properties']) - posted} already present") + + +def main(): + """Command line: validate the run documents, then upload them.""" + ap = argparse.ArgumentParser(description=__doc__, formatter_class=argparse.RawDescriptionHelpFormatter) + ap.add_argument("documents", nargs="+", metavar="RUN.JSON", help="run documents, as a parser writes them") + ap.add_argument("--host", default=os.environ.get("MAT3RA_HOST", "localhost:3000"), + help="web app host or URL (or MAT3RA_HOST); https unless localhost, e.g. dev.mat3ra.com") + ap.add_argument("--account", help="slug of the account the data belongs to (reads scoped to it, writes owned by it)") + ap.add_argument("--files", nargs="*", default=["records"], metavar="GROUP", + help=f"which file groups to upload ({', '.join(FILE_GROUPS)}); pass --files with no value " + "to upload none and keep the raw records in each measurement's metadata") + ap.add_argument("--no-properties", action="store_true", + help="upload the run without its properties - the files still carry everything they were derived from") + ap.add_argument("--dry-run", action="store_true", help="validate the documents and stop") + a = ap.parse_args() + + runs = [load(path) for path in a.documents] + errors = sum(validate(run) for run in runs) + print("validation:", "OK" if errors == 0 else f"{errors} invalid documents") + if errors or a.dry_run: + sys.exit(1 if errors else 0) + + url = urllib.parse.urlsplit(base_url(a.host)) + address = {"host": url.hostname, "port": url.port or (443 if url.scheme == "https" else 80), "secure": url.scheme == "https"} + client = APIClient.authenticate(**address) + if a.account: + client = APIClient.authenticate(account_id=account_id(client, a.account), **address) + for run in runs: + upload(client, run, files=a.files, properties=not a.no_properties) + + +if __name__ == "__main__": + main() diff --git a/examples/measurement/upload_spm_run.ipynb b/examples/measurement/upload_spm_run.ipynb new file mode 100644 index 000000000..790fd4eab --- /dev/null +++ b/examples/measurement/upload_spm_run.ipynb @@ -0,0 +1,225 @@ +{ + "cells": [ + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "# Overview\n", + "\n", + "This example uploads one scanning probe microscopy run to the platform: the run folder becomes a Sample Set with one Sample per measured position, a Measurement Set with one Measurement per Sample and its Setup, the run's records as files, and one hysteresis-loop Property per Sample.\n", + "A run folder is what the instrument exports — `summary.json` with the recipe, the session and one record per measured point, and `loops/` with the raw curves — and re-running the notebook adds only what is missing." + ] + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "## Install the API client\n", + "\n", + "The samples, measurements and files endpoints are not released yet, so the client is installed from its branch until it merges. Restart the kernel after this cell." + ] + }, + { + "cell_type": "code", + "metadata": {}, + "source": [ + "\n", + "%pip install -q -r requirements.txt" + ], + "outputs": [], + "execution_count": null + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "## Set Parameters\n", + "\n", + "- **HOST**: platform the run is uploaded to\n", + "- **RUN_DIR**: the run folder beside this notebook — what the instrument exports, with `summary.json` and `loops/` inside it\n", + "- **PHYSICAL_ID**: the identifier written on the physical piece the measured positions are part of — every Sample carries it\n", + "- **FILES**: which files to upload per measurement" + ] + }, + { + "cell_type": "code", + "metadata": {}, + "source": [ + "import urllib.parse\n", + "\n", + "HOST = \"https://alphafilm.mat3ra.com\"\n", + "\n", + "# NOTE: generate at https://alphafilm.mat3ra.com/demo/preferences API Tokens\n", + "ACCOUNT_ID = \"pR5gsAQJrJYJuapM3\"\n", + "AUTH_TOKEN = \"gFEku1YpH4gTF57g0nYUj2C_UYM5brGSZXjSoCywI-H\"\n", + "\n", + "RUN_DIR = \"\" # folder with summary.json, records etc\n", + "PHYSICAL_ID = \"\" # the physical wafer's identifier (e.g. PDAC_COM5_01448)\n", + "FILES = [\"records\", \"loops\"] # Which folders to upload as files per measurement.\n", + "\n", + "url = urllib.parse.urlsplit(HOST)\n", + "address = {\n", + " \"host\": url.hostname,\n", + " \"port\": url.port or (443 if url.scheme == \"https\" else 80),\n", + " \"secure\": url.scheme == \"https\",\n", + "}" + ], + "outputs": [], + "execution_count": null + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "## Authenticate and initialize API client\n", + "\n", + "### Authenticate\n", + "Authenticate in the browser (OIDC device flow) or via JupyterLite host injection. Credentials are stored in environment variables.\n", + "\n", + "### Initialize API client\n", + "Create an authenticated API client and resolve the owner account ID." + ] + }, + { + "cell_type": "code", + "metadata": {}, + "source": [ + "from mat3ra.notebooks_utils.auth import authenticate\n", + "import os\n", + "\n", + "os.environ[\"ACCOUNT_ID\"] = ACCOUNT_ID\n", + "os.environ[\"AUTH_TOKEN\"] = AUTH_TOKEN\n", + "\n", + "\n", + "# NOTE: uncomment to login with OIDC interactively\n", + "# os.environ[\"API_HOST\"] = address[\"host\"]\n", + "# os.environ[\"API_PORT\"] = str(address[\"port\"])\n", + "# os.environ[\"API_SECURE\"] = str(address[\"secure\"])\n", + "# await authenticate()" + ], + "outputs": [], + "execution_count": null + }, + { + "cell_type": "code", + "metadata": {}, + "source": [ + "from mat3ra.api_client import APIClient\n", + "\n", + "client = APIClient.authenticate(**address)" + ], + "outputs": [], + "execution_count": null + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "# Imports" + ] + }, + { + "cell_type": "code", + "metadata": {}, + "source": [ + "from pathlib import Path\n", + "\n", + "from parse_utk import parse\n", + "from run_document import load, serialize\n", + "from upload_run import account_id, upload" + ], + "outputs": [], + "execution_count": null + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "## Parse the run folder\n", + "\n", + "Read the run folder into the documents the platform stores. Nothing is uploaded yet." + ] + }, + { + "cell_type": "code", + "metadata": {}, + "source": [ + "# reading the run folder and writing the run document: nothing here talks to the platform\n", + "parsed = parse(Path(RUN_DIR), PHYSICAL_ID)\n", + "document_path = serialize(parsed, \"parsed\")\n", + "run = load(document_path)\n", + "file_count = sum(len(files) for files in run[\"files\"].values())\n", + "print(\n", + " f\"{run['physicalId']}: {len(run['samples'])} samples (ordered set) · run {run['run']}: \"\n", + " f\"{len(run['measurements'])} measurements (ordered set, one per sample) · {len(run['set_files'])} set files \"\n", + " f\"-> {file_count} measurement files · {len(run['properties'])} samples with a combined loop \"\n", + " f\"· no curves: {len(run['skipped'])} samples\"\n", + ")\n", + "print(\"run document:\", document_path)" + ], + "outputs": [], + "execution_count": null + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "## Upload the run\n", + "\n", + "Create the Sample Set and its Samples, the Measurement Set and one Measurement per Sample, the files and the loop Properties." + ] + }, + { + "cell_type": "code", + "metadata": {}, + "source": [ + "upload(client, run, files=FILES)" + ], + "outputs": [], + "execution_count": null + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "## Find the run in the web app\n", + "\n", + "The run is a folder in the account's Measurements tab, named after the run." + ] + }, + { + "cell_type": "code", + "metadata": {}, + "source": [ + "print(f\"Open {HOST}, your account's Measurements tab: {run['run']}\")" + ], + "outputs": [], + "execution_count": null + } + ], + "metadata": { + "colab": { + "name": "upload_spm_run.ipynb", + "provenance": [] + }, + "kernelspec": { + "display_name": "Python 3", + "language": "python", + "name": "python3" + }, + "language_info": { + "codemirror_mode": { + "name": "ipython", + "version": 3 + }, + "file_extension": ".py", + "mimetype": "text/x-python", + "name": "python", + "nbconvert_exporter": "python", + "pygments_lexer": "ipython3", + "version": "3.8.6" + } + }, + "nbformat": 4, + "nbformat_minor": 1 +}