From 8d3dfb49ecf28219e719ca49ef9bb7eb9f104338 Mon Sep 17 00:00:00 2001 From: Biowilko Date: Mon, 24 Aug 2026 12:05:54 +0100 Subject: [PATCH 1/3] Don't parse bucket names anymore --- roz_scripts/general/s3_controller.py | 6 +- roz_scripts/general/s3_matcher.py | 77 +++-- roz_scripts/general/s3_onyx_updates.py | 32 +-- roz_scripts/utils/config.py | 222 +++++++++++++++ tests/fixtures/test_config.json | 1 + tests/fixtures/test_config_adversarial.json | 68 +++++ tests/test_config.py | 300 ++++++++++++++++++++ tests/test_csv_update.py | 16 +- tests/test_s3_matcher.py | 4 +- 9 files changed, 677 insertions(+), 49 deletions(-) create mode 100644 tests/fixtures/test_config_adversarial.json create mode 100644 tests/test_config.py diff --git a/roz_scripts/general/s3_controller.py b/roz_scripts/general/s3_controller.py index 8ff9912..cbb9cf0 100644 --- a/roz_scripts/general/s3_controller.py +++ b/roz_scripts/general/s3_controller.py @@ -8,6 +8,8 @@ import copy import requests +from roz_scripts.utils.config import TEST_FLAGS + S3_ENDPOINT_URL = "https://s3.climb.ac.uk" REQUESTS_TIMEOUT = 30 @@ -180,7 +182,7 @@ def create_config_map(config_dict: dict) -> dict: desired_labels = re.findall(r"{(\w*)}", bucket_config["name_layout"]) for platform in config["file_specs"].keys(): - for test_flag in ["prod", "test"]: + for test_flag in TEST_FLAGS: try: namespace = {} @@ -203,7 +205,7 @@ def create_config_map(config_dict: dict) -> dict: desired_labels = re.findall(r"{(\w*)}", bucket_config["name_layout"]) for platform in config["file_specs"].keys(): - for test_flag in ["prod", "test"]: + for test_flag in TEST_FLAGS: try: namespace = {} diff --git a/roz_scripts/general/s3_matcher.py b/roz_scripts/general/s3_matcher.py index 4bc83ca..2ae11c3 100644 --- a/roz_scripts/general/s3_matcher.py +++ b/roz_scripts/general/s3_matcher.py @@ -7,12 +7,14 @@ ) from roz_scripts.utils.health import HealthState, get_health_dir from roz_scripts.general.s3_controller import create_config_map +from roz_scripts.utils.config import load_config, parse_ingest_bucket_name, ConfigError from varys import Varys import boto3 from botocore.client import BaseClient from botocore.exceptions import ClientError +import logging import uuid import time import json @@ -125,12 +127,19 @@ def gen_s3_uri(bucket_name: str, key: str) -> str: return f"s3://{bucket_name}/{key}" -def parse_existing_objects(existing_objects: dict, config_dict: dict) -> dict: +def parse_existing_objects( + existing_objects: dict, config_dict: dict, log: logging.Logger | None = None +) -> dict: """Parses existing objects into a dictionary of artifacts. Args: existing_objects (dict): Dictionary of existing objects from func get_existing_objects config_dict (dict): Dictionary containing the config file + log (logging.Logger | None): Logger object, for reporting a bucket + name that doesn't resolve against config_dict. Buckets reaching + this point were themselves discovered from config_dict, so this + should never happen in practice - it's a defensive skip, not an + expected path. Returns: dict: Dictionary of artifacts @@ -139,12 +148,18 @@ def parse_existing_objects(existing_objects: dict, config_dict: dict) -> dict: parsed_objects = {} for bucket_name, objs in existing_objects.items(): - project, site_str, platform, test_flag = bucket_name.split("-") - - if "." in site_str: - site = site_str.split(".")[-2] - else: - site = site_str + try: + parsed_bucket_name = parse_ingest_bucket_name(config_dict, bucket_name) + except ConfigError as e: + if log is not None: + log.error(f"Skipping unresolvable bucket {bucket_name!r}: {e}") + continue + + project = parsed_bucket_name["project"] + site_str = parsed_bucket_name["raw_site"] + site = parsed_bucket_name["site"] + platform = parsed_bucket_name["platform"] + test_flag = parsed_bucket_name["test_flag"] for obj in objs: # Ignore test key, s3_controller uses it to check if the bucket is correctly configured @@ -227,7 +242,7 @@ def is_artifact_dict_complete( def parse_new_object_message( existing_object_dict: dict, new_object_message: dict, config_dict: dict -) -> tuple[bool, dict, tuple, dict]: +) -> tuple[bool, dict, tuple | None, dict | None]: """Parses a new object message and adds it to the existing object dict. Args: @@ -236,7 +251,13 @@ def parse_new_object_message( config_dict (dict): Dictionary parsed from the config file Returns: - tuple[bool, dict, tuple, dict]: Tuple containing a boolean indicating if the artifact is complete, the updated existing object dict, the index tuple, and the parsed bucket name + tuple[bool, dict, tuple | None, dict | None]: Tuple containing a + boolean indicating if the artifact is complete, the updated + existing object dict, the index tuple, and the parsed bucket + name. index_tuple and the parsed bucket name are both None if + the message's bucket name doesn't resolve against config_dict at + all (e.g. a stale/foreign bucket) - the caller must check for + this before unpacking index_tuple. """ # There should only ever be one record here @@ -244,19 +265,12 @@ def parse_new_object_message( bucket_name = record["s3"]["bucket"]["name"] - parsed_bucket_name = { - x: y - for x, y in zip( - ("project", "site_str", "platform", "test_flag"), bucket_name.split("-") - ) - } - - # project, site_str, platform, test_flag = parsed_bucket_name + try: + parsed_bucket_name = parse_ingest_bucket_name(config_dict, bucket_name) + except ConfigError: + return (False, existing_object_dict, None, None) - if "." in parsed_bucket_name["site_str"]: - site = parsed_bucket_name["site_str"].split(".")[-2] - else: - site = parsed_bucket_name["site_str"] + site = parsed_bucket_name["site"] object_key = record["s3"]["object"]["key"] @@ -332,7 +346,7 @@ def parse_new_object_message( existing_object_dict[index_tuple]["objects"][extension] = record if extension == ".csv": - existing_object_dict[index_tuple]["raw_site"] = parsed_bucket_name["site_str"] + existing_object_dict[index_tuple]["raw_site"] = parsed_bucket_name["raw_site"] return ( is_artifact_dict_complete( @@ -419,8 +433,11 @@ def main(): log_level=os.environ["INGEST_LOG_LEVEL"], ) - with open(os.environ["ROZ_CONFIG_JSON"], "r") as f: - config_dict = json.load(f) + try: + config_dict = load_config() + except ConfigError as e: + print(f"Invalid roz config: {e}", file=sys.stderr) + sys.exit(3) config_map = create_config_map(config_dict=config_dict) @@ -439,7 +456,7 @@ def main(): objects = get_existing_objects(s3_client=s3_client, to_check=buckets) existing_object_dict = parse_existing_objects( - existing_objects=objects, config_dict=config_dict + existing_objects=objects, config_dict=config_dict, log=log ) health = HealthState(get_health_dir()) @@ -470,6 +487,16 @@ def main(): ) ) + if index_tuple is None: + bucket_name = message_dict["Records"][0]["s3"]["bucket"]["name"] + failure_message = f"Bucket name {bucket_name!r} does not match any known bucket layout, skipping message" + log.error(failure_message) + send_admin_alert( + varys_client, source="s3_matcher", description=failure_message + ) + continue + + assert parsed_bucket_name is not None artifact, project, site, platform, test_flag = index_tuple if not artifact: diff --git a/roz_scripts/general/s3_onyx_updates.py b/roz_scripts/general/s3_onyx_updates.py index 469d5ee..caf3852 100644 --- a/roz_scripts/general/s3_onyx_updates.py +++ b/roz_scripts/general/s3_onyx_updates.py @@ -17,6 +17,12 @@ ) from roz_scripts.general.s3_matcher import parse_object_key from roz_scripts.utils.health import HealthState, get_health_dir +from roz_scripts.utils.config import ( + load_config, + parse_ingest_bucket_name, + site_bucket, + ConfigError, +) from varys import Varys from onyx import ( @@ -182,18 +188,13 @@ def csv_update(parsed_message, config_dict, log): record = parsed_message["Records"][0] bucket_name = record["s3"]["bucket"]["name"] - parsed_bucket_name = { - x: y - for x, y in zip( - ("project", "site_str", "platform", "test_flag"), - bucket_name.split("-"), - ) - } + try: + parsed_bucket_name = parse_ingest_bucket_name(config_dict, bucket_name) + except ConfigError as e: + log.error(f"Skipping unresolvable bucket {bucket_name!r}: {e}") + return (True, False) - if "." in parsed_bucket_name["site_str"]: - site = parsed_bucket_name["site_str"].split(".")[-2] - else: - site = parsed_bucket_name["site_str"] + site = parsed_bucket_name["site"] # ignore files from test buckets if parsed_bucket_name["test_flag"] == "test": @@ -283,7 +284,7 @@ def csv_update(parsed_message, config_dict, log): "run_index": parsed_object_key["run_index"], "run_id": parsed_object_key["run_id"], "site": site, - "site_str": parsed_bucket_name["site_str"], + "site_str": parsed_bucket_name["raw_site"], "files": { ".csv": { "uri": f"s3://{bucket_name}/{record['s3']['object']['key']}", @@ -314,7 +315,7 @@ def csv_update(parsed_message, config_dict, log): payload["update_status"] = "failed" s3_client.put_object( - Bucket=f"{parsed_bucket_name['project']}-{parsed_bucket_name['site_str']}-results", + Bucket=site_bucket(config_dict, parsed_bucket_name["project"], parsed_bucket_name["raw_site"], "results"), Key=f"{payload['artifact']}.update.json", Body=json.dumps(payload), ) @@ -358,7 +359,7 @@ def csv_update(parsed_message, config_dict, log): payload["update_status"] = "failed" if update_failure else "success" s3_client.put_object( - Bucket=f"{parsed_bucket_name['project']}-{parsed_bucket_name['site_str']}-results", + Bucket=site_bucket(config_dict, parsed_bucket_name["project"], parsed_bucket_name["raw_site"], "results"), Key=f"{payload['artifact']}.update.json", Body=json.dumps(payload), ) @@ -378,8 +379,7 @@ def run(args): auto_acknowledge=False, ) - with open(os.environ["ROZ_CONFIG_JSON"], "r") as f: - config_dict = json.load(f) + config_dict = load_config() health = HealthState(get_health_dir()) diff --git a/roz_scripts/utils/config.py b/roz_scripts/utils/config.py index 20cf253..6214771 100644 --- a/roz_scripts/utils/config.py +++ b/roz_scripts/utils/config.py @@ -1,3 +1,5 @@ +from dataclasses import dataclass, field +import itertools import json import os import regex as re @@ -7,6 +9,9 @@ class ConfigError(ValueError): """Raised when the roz config can't be loaded, or a bucket lookup fails against it""" +TEST_FLAGS = ("prod", "test") + + _config_cache: dict[str, dict] = {} @@ -143,3 +148,220 @@ def project_bucket_uri(config: dict, project: str, bucket: str, key: str, **kwar str: An "s3://bucket/key" URI """ return f"s3://{project_bucket(config, project, bucket, **kwargs)}/{key}" + + +@dataclass(frozen=True) +class BucketMatch: + """The result of resolving a bucket name back to the config entry that produces it + + Attributes: + bucket_name (str): The bucket name that was parsed + project (str): The project the bucket belongs to, e.g. "mscape" + bucket (str): The project_buckets/site_buckets key, e.g. "ingest" + scope (str): Either "project" or "site" + site (str | None): The exact config `sites` key, for scope="site". None for scope="project" + platform (str | None): The platform, if the bucket's name_layout includes {platform} + test_flag (str | None): "prod" or "test", if the bucket's name_layout includes {test_flag} + fields (dict): Every placeholder value the name_layout was expanded with + """ + + bucket_name: str + project: str + bucket: str + scope: str + site: str | None + platform: str | None + test_flag: str | None + fields: dict = field(default_factory=dict) + + +def short_site(site: str) -> str: + """Derive the short form of a site name used for exchange names and dedup keys + + e.g. "gpha.ukhsa.mscape" -> "ukhsa". Note this is a lossy heuristic that + can collide between distinct sites sharing a second-to-last dotted + component - it exists here, centralised, purely to match pre-existing + behaviour rather than as an endorsement of the scheme. + + Args: + site (str): The full config `sites` key + + Returns: + str: The short site name + """ + return site.split(".")[-2] if "." in site else site + + +def _label_domain(label: str, project: str, project_config: dict, site: str | None): + """Return every value a name_layout placeholder can take, or None if it can't be enumerated""" + if label == "project": + return [project] + if label == "site": + return [site] if site is not None else None + if label == "platform": + return list(project_config.get("file_specs", {}).keys()) + if label == "test_flag": + return list(TEST_FLAGS) + return None + + +def _index_layout( + index: dict, project: str, project_config: dict, bucket: str, bucket_config: dict, scope: str, site: str | None +) -> None: + name_layout = bucket_config["name_layout"] + labels = re.findall(r"{(\w*)}", name_layout) + + domains = [] + for label in labels: + domain = _label_domain(label, project, project_config, site) + if domain is None: + raise ConfigError( + f"Cannot enumerate values for placeholder '{label}' in bucket layout " + f"{name_layout!r} (project={project!r}, bucket={bucket!r})" + ) + domains.append(domain) + + for combo in itertools.product(*domains) if domains else [()]: + fields = dict(zip(labels, combo)) + bucket_name = name_layout.format(**fields) + + match = BucketMatch( + bucket_name=bucket_name, + project=project, + bucket=bucket, + scope=scope, + site=fields.get("site"), + platform=fields.get("platform"), + test_flag=fields.get("test_flag"), + fields=fields, + ) + + if bucket_name in index: + existing = index[bucket_name] + raise ConfigError( + f"Bucket name {bucket_name!r} is ambiguous: matches both " + f"(project={existing.project!r}, bucket={existing.bucket!r}, scope={existing.scope!r}) " + f"and (project={project!r}, bucket={bucket!r}, scope={scope!r})" + ) + + index[bucket_name] = match + + +def _build_bucket_index(config: dict) -> dict: + index: dict = {} + + for project, project_config in config.get("configs", {}).items(): + for bucket, bucket_config in project_config.get("project_buckets", {}).items(): + _index_layout(index, project, project_config, bucket, bucket_config, scope="project", site=None) + + for site in project_config.get("sites", {}): + for bucket, bucket_config in project_config.get("site_buckets", {}).items(): + _index_layout(index, project, project_config, bucket, bucket_config, scope="site", site=site) + + return index + + +# Keyed by id(config) rather than the config dict itself (which isn't +# hashable in general) - safe because load_config() caches and returns the +# same long-lived dict object for a given path, which is how every caller +# obtains a config in practice. +_bucket_index_cache: dict[int, dict] = {} + + +def _get_bucket_index(config: dict) -> dict: + key = id(config) + if key not in _bucket_index_cache: + _bucket_index_cache[key] = _build_bucket_index(config) + return _bucket_index_cache[key] + + +def parse_bucket_name( + config: dict, bucket_name: str, *, bucket: str | None = None, scope: str | None = None +) -> BucketMatch: + """Resolve a bucket name back to the config entry that produces it + + This is the inverse of project_bucket()/site_bucket(): rather than + inferring structure from the bucket name's characters (fragile - see the + positional-split parsing this replaces), it enumerates every bucket name + the config can produce and looks the input up in that index. This means + a bucket name can only ever resolve to a real, currently-configured + bucket, and an ambiguous config (two layouts producing the same name) is + caught at index-build time rather than on a live notification. + + Args: + config (dict): The loaded roz config, from load_config() + bucket_name (str): The bucket name to resolve + bucket (str | None): If given, require the match's bucket key to equal this + scope (str | None): If given, require the match's scope ("project" or "site") to equal this + + Returns: + BucketMatch: The resolved match + + Raises: + ConfigError: If the bucket name doesn't match any known bucket, or + doesn't satisfy the `bucket`/`scope` filters + """ + index = _get_bucket_index(config) + + match = index.get(bucket_name) + if match is None: + raise ConfigError(f"Bucket name {bucket_name!r} does not match any known bucket layout") + + if bucket is not None and match.bucket != bucket: + raise ConfigError( + f"Bucket name {bucket_name!r} resolved to bucket {match.bucket!r}, expected {bucket!r}" + ) + + if scope is not None and match.scope != scope: + raise ConfigError( + f"Bucket name {bucket_name!r} resolved to scope {match.scope!r}, expected {scope!r}" + ) + + return match + + +def try_parse_bucket_name(config: dict, bucket_name: str, **kwargs) -> BucketMatch | None: + """Like parse_bucket_name(), but returns None instead of raising on no/bad match + + Args: + config (dict): The loaded roz config, from load_config() + bucket_name (str): The bucket name to resolve + **kwargs: Passed through to parse_bucket_name() (bucket, scope) + + Returns: + BucketMatch | None: The resolved match, or None + """ + try: + return parse_bucket_name(config, bucket_name, **kwargs) + except ConfigError: + return None + + +def parse_ingest_bucket_name(config: dict, bucket_name: str) -> dict: + """Resolve an ingest bucket name into the fields s3_matcher/s3_onyx_updates need + + Args: + config (dict): The loaded roz config, from load_config() + bucket_name (str): The ingest bucket name to resolve + + Returns: + dict: {"project", "raw_site", "site", "platform", "test_flag", "scope"}. + For a project-level ingest bucket (no real site involved), + raw_site/site are "public", matching this codebase's pre-existing + convention for that case. + + Raises: + ConfigError: If the bucket name doesn't resolve to an "ingest" bucket + """ + match = parse_bucket_name(config, bucket_name, bucket="ingest") + + raw_site = match.site or "public" + + return { + "project": match.project, + "raw_site": raw_site, + "site": short_site(raw_site), + "platform": match.platform, + "test_flag": match.test_flag, + "scope": match.scope, + } diff --git a/tests/fixtures/test_config.json b/tests/fixtures/test_config.json index df5a798..c3b6eb7 100644 --- a/tests/fixtures/test_config.json +++ b/tests/fixtures/test_config.json @@ -6,6 +6,7 @@ "artifact_layout": "project|run_index|run_id", "sites": { "birm.mscape": "analysis", + "subteam1.birm.mscape": "analysis", "site1.mscape": "analysis", "site2.mscape": "uploader", "site3.mscape": "uploader", diff --git a/tests/fixtures/test_config_adversarial.json b/tests/fixtures/test_config_adversarial.json new file mode 100644 index 0000000..0b1dc14 --- /dev/null +++ b/tests/fixtures/test_config_adversarial.json @@ -0,0 +1,68 @@ +{ + "version": 1, + "pathogen_configs": ["adversarialscape"], + "configs": { + "adversarialscape": { + "artifact_layout": "project|run_index|run_id", + "sites": { + "adversarialscape": "analysis", + "gpha.ukhsa.adversarialscape": "analysis", + "tarzet.ukhsa.adversarialscape": "uploader", + "hyphen-site.adversarialscape": "uploader", + "collide.adversarialscape": "uploader", + "collide": "uploader" + }, + "csv_updates": false, + "bucket_policies": { + "site_ingest": ["get", "put", "list", "delete"], + "project_ingest": ["get", "put", "list", "delete"], + "project_read": ["get"] + }, + "site_buckets": { + "ingest": { + "name_layout": "{project}-{site}-{platform}-{test_flag}", + "policy": {"analysis": "site_ingest", "uploader": "site_ingest"}, + "owner": "{site}" + }, + "results": { + "name_layout": "{project}-{site}-fakeresults", + "policy": {"analysis": "site_ingest", "uploader": "site_ingest"}, + "owner": "{site}" + } + }, + "notification_bucket_configs": { + "ingest": { + "rmq_exchange": "inbound-s3", + "rmq_queue_env": "s3_matcher", + "amqps": false + } + }, + "project_buckets": { + "ingest": { + "name_layout": "{project}-public-{platform}-{test_flag}", + "policy": {"analysis": "project_ingest"}, + "owner": "admin" + }, + "published_reads": { + "name_layout": "{project}-fake-adversarial-reads", + "policy": {"analysis": "project_read"}, + "owner": "admin" + } + }, + "file_specs": { + "illumina": { + ".1.fastq.gz": {"layout": "project.run_index.run_id.direction.ftype.gzip"}, + ".2.fastq.gz": {"layout": "project.run_index.run_id.direction.ftype.gzip"}, + ".csv": {"layout": "project.run_index.run_id.ftype"} + }, + "illumina.se": { + ".fastq.gz": {"layout": "project.run_index.run_id.ftype.gzip"}, + ".csv": {"layout": "project.run_index.run_id.ftype"} + }, + "noplatform": { + ".csv": {"layout": "project.run_index.run_id.ftype"} + } + } + } + } +} diff --git a/tests/test_config.py b/tests/test_config.py new file mode 100644 index 0000000..a1a6ac5 --- /dev/null +++ b/tests/test_config.py @@ -0,0 +1,300 @@ +import itertools +import os +import unittest + +from roz_scripts.utils.config import ( + ConfigError, + TEST_FLAGS, + load_config, + project_bucket, + site_bucket, + parse_bucket_name, + try_parse_bucket_name, + parse_ingest_bucket_name, + short_site, +) + +DIR = os.path.dirname(__file__) +TEST_CONFIG_PATH = os.path.join(DIR, "fixtures", "test_config.json") +ADVERSARIAL_CONFIG_PATH = os.path.join(DIR, "fixtures", "test_config_adversarial.json") + + +class test_short_site(unittest.TestCase): + def test_dotted_site_returns_second_to_last_component(self): + self.assertEqual(short_site("gpha.ukhsa.mscape"), "ukhsa") + + def test_bare_site_returns_itself(self): + self.assertEqual(short_site("mscape"), "mscape") + + def test_two_component_site(self): + self.assertEqual(short_site("birm.mscape"), "birm") + + +class test_parse_bucket_name_round_trip(unittest.TestCase): + """The single highest-value test here: for every project/scope/bucket/ + site/platform/test_flag combination the real fixture config can produce, + construct the bucket name and assert parsing it recovers the same + coordinates. If this ever fails, the constructor and parser have + diverged - which is exactly the class of bug this whole change exists + to make structurally impossible. + """ + + def setUp(self): + self.config = load_config(TEST_CONFIG_PATH) + + def test_all_project_buckets_round_trip(self): + checked = 0 + for project, project_config in self.config["configs"].items(): + for bucket, bucket_config in project_config["project_buckets"].items(): + labels = set(_placeholders(bucket_config["name_layout"])) + for platform, test_flag in _platform_test_flag_combos( + project_config, labels + ): + kwargs = {} + if "platform" in labels: + kwargs["platform"] = platform + if "test_flag" in labels: + kwargs["test_flag"] = test_flag + + bucket_name = project_bucket(self.config, project, bucket, **kwargs) + match = parse_bucket_name(self.config, bucket_name) + + self.assertEqual(match.project, project) + self.assertEqual(match.bucket, bucket) + self.assertEqual(match.scope, "project") + self.assertIsNone(match.site) + checked += 1 + + self.assertGreater(checked, 0) + + def test_all_site_buckets_round_trip(self): + checked = 0 + for project, project_config in self.config["configs"].items(): + for site in project_config["sites"]: + for bucket, bucket_config in project_config["site_buckets"].items(): + labels = set(_placeholders(bucket_config["name_layout"])) + for platform, test_flag in _platform_test_flag_combos( + project_config, labels + ): + kwargs = {} + if "platform" in labels: + kwargs["platform"] = platform + if "test_flag" in labels: + kwargs["test_flag"] = test_flag + + bucket_name = site_bucket( + self.config, project, site, bucket, **kwargs + ) + match = parse_bucket_name(self.config, bucket_name) + + self.assertEqual(match.project, project) + self.assertEqual(match.bucket, bucket) + self.assertEqual(match.scope, "site") + self.assertEqual(match.site, site) + checked += 1 + + self.assertGreater(checked, 0) + + +def _placeholders(name_layout): + import regex as re + + return re.findall(r"{(\w*)}", name_layout) + + +def _platform_test_flag_combos(project_config, labels): + platforms = project_config["file_specs"].keys() if "platform" in labels else [None] + test_flags = TEST_FLAGS if "test_flag" in labels else [None] + return itertools.product(platforms, test_flags) + + +class test_parse_bucket_name_adversarial(unittest.TestCase): + """Cases the tame fixture can't exercise: hyphenated sites, multi-level + dotted sites, short-site collisions, and the site/project ingest layout + collision-by-design (site named "public" would collide with the + project-level pseudo-site, which is exactly why construct-and-index + - not a free-form regex - is the right approach). + """ + + def setUp(self): + self.config = load_config(ADVERSARIAL_CONFIG_PATH) + + def test_hyphenated_site_parses_correctly(self): + # This is the exact case the old `bucket_name.split("-")` positional + # parse cannot handle - a site name that itself contains a hyphen + # produces more than 4 fields. + bucket_name = "adversarialscape-hyphen-site.adversarialscape-illumina-prod" + match = parse_bucket_name(self.config, bucket_name) + + self.assertEqual(match.project, "adversarialscape") + self.assertEqual(match.bucket, "ingest") + self.assertEqual(match.scope, "site") + self.assertEqual(match.site, "hyphen-site.adversarialscape") + self.assertEqual(match.platform, "illumina") + self.assertEqual(match.test_flag, "prod") + + def test_multi_level_dotted_site(self): + bucket_name = "adversarialscape-gpha.ukhsa.adversarialscape-illumina-prod" + match = parse_bucket_name(self.config, bucket_name) + self.assertEqual(match.site, "gpha.ukhsa.adversarialscape") + + def test_illumina_se_platform_not_confused_with_illumina(self): + bucket_name = "adversarialscape-adversarialscape-illumina.se-prod" + match = parse_bucket_name(self.config, bucket_name) + self.assertEqual(match.platform, "illumina.se") + + def test_bare_project_site(self): + bucket_name = "adversarialscape-adversarialscape-illumina-prod" + match = parse_bucket_name(self.config, bucket_name) + self.assertEqual(match.site, "adversarialscape") + + def test_project_scope_ingest_bucket_uses_public_pseudo_site(self): + bucket_name = "adversarialscape-public-illumina-prod" + match = parse_bucket_name(self.config, bucket_name) + + self.assertEqual(match.scope, "project") + self.assertIsNone(match.site) + + parsed = parse_ingest_bucket_name(self.config, bucket_name) + self.assertEqual(parsed["raw_site"], "public") + self.assertEqual(parsed["site"], "public") + + def test_short_site_collision_does_not_cause_bucket_name_ambiguity(self): + # "collide" and "collide.adversarialscape" share a short site name, + # but remain distinct full site names and therefore distinct bucket + # names - no ConfigError building the index over this fixture. + m1 = parse_bucket_name( + self.config, "adversarialscape-collide-illumina-prod" + ) + m2 = parse_bucket_name( + self.config, "adversarialscape-collide.adversarialscape-illumina-prod" + ) + self.assertEqual(m1.site, "collide") + self.assertEqual(m2.site, "collide.adversarialscape") + assert m1.site is not None and m2.site is not None + self.assertEqual(short_site(m1.site), short_site(m2.site)) + + +class test_parse_bucket_name_malformed(unittest.TestCase): + def setUp(self): + self.config = load_config(TEST_CONFIG_PATH) + + def _assert_rejected(self, bucket_name): + with self.assertRaises(ConfigError): + parse_bucket_name(self.config, bucket_name) + self.assertIsNone(try_parse_bucket_name(self.config, bucket_name)) + + def test_empty_string(self): + self._assert_rejected("") + + def test_too_few_fields(self): + self._assert_rejected("mscape") + self._assert_rejected("mscape-birm.mscape") + self._assert_rejected("mscape-birm.mscape-illumina") + + def test_too_many_fields(self): + self._assert_rejected("mscape-birm.mscape-illumina-prod-extra") + + def test_invalid_test_flag(self): + self._assert_rejected("mscape-birm.mscape-illumina-staging") + + def test_unknown_platform(self): + self._assert_rejected("mscape-birm.mscape-nanopore-prod") + + def test_unknown_site(self): + self._assert_rejected("mscape-nonexistent.mscape-illumina-prod") + + def test_unknown_project(self): + self._assert_rejected("notaproject-birm.mscape-illumina-prod") + + def test_case_sensitivity(self): + self._assert_rejected("MSCAPE-BIRM.MSCAPE-ILLUMINA-PROD") + + def test_empty_field(self): + self._assert_rejected("mscape--illumina-prod") + + def test_unescaped_dot_regression(self): + # A hand-rolled regex inverse (`{site}` as `.+` etc) could easily let + # a "." in a literal position match any character. Construct-and- + # index has no such risk, but pin it as a regression guard anyway. + self._assert_rejected("mscapeXbirmXmscape-illumina-prod") + + +class test_parse_bucket_name_filters(unittest.TestCase): + def setUp(self): + self.config = load_config(TEST_CONFIG_PATH) + + def test_bucket_filter_accepts_matching(self): + bucket_name = "mscape-birm.mscape-ont-prod" + match = parse_bucket_name(self.config, bucket_name, bucket="ingest") + self.assertEqual(match.bucket, "ingest") + + def test_bucket_filter_rejects_non_matching(self): + bucket_name = project_bucket(self.config, "mscape", "published_reads") + with self.assertRaises(ConfigError): + parse_bucket_name(self.config, bucket_name, bucket="ingest") + + def test_scope_filter_rejects_non_matching(self): + bucket_name = site_bucket(self.config, "mscape", "birm.mscape", "results") + with self.assertRaises(ConfigError): + parse_bucket_name(self.config, bucket_name, scope="project") + + +class test_build_bucket_index_collision_detection(unittest.TestCase): + def test_colliding_layouts_raise_configerror_at_build_time(self): + from roz_scripts.utils.config import _build_bucket_index + + colliding_config = { + "configs": { + "proj": { + "sites": {"siteA": "analysis"}, + "file_specs": {"ont": {}}, + "site_buckets": { + "ingest": { + "name_layout": "{project}-{site}", + "policy": {}, + "owner": "{site}", + } + }, + "project_buckets": { + "collider": { + "name_layout": "{project}-siteA", + "policy": {}, + "owner": "admin", + } + }, + } + } + } + + with self.assertRaises(ConfigError): + _build_bucket_index(colliding_config) + + +class test_parse_ingest_bucket_name(unittest.TestCase): + def setUp(self): + self.config = load_config(TEST_CONFIG_PATH) + + def test_returns_expected_keys(self): + bucket_name = "mscape-birm.mscape-ont-prod" + parsed = parse_ingest_bucket_name(self.config, bucket_name) + + self.assertEqual( + set(parsed.keys()), + {"project", "raw_site", "site", "platform", "test_flag", "scope"}, + ) + self.assertEqual(parsed["project"], "mscape") + self.assertEqual(parsed["raw_site"], "birm.mscape") + self.assertEqual(parsed["site"], "birm") + self.assertEqual(parsed["platform"], "ont") + self.assertEqual(parsed["test_flag"], "prod") + self.assertEqual(parsed["scope"], "site") + + def test_rejects_non_ingest_bucket(self): + bucket_name = project_bucket(self.config, "mscape", "published_reads") + with self.assertRaises(ConfigError): + parse_ingest_bucket_name(self.config, bucket_name) + + +if __name__ == "__main__": + unittest.main() diff --git a/tests/test_csv_update.py b/tests/test_csv_update.py index ab2d4cf..6b1f20e 100644 --- a/tests/test_csv_update.py +++ b/tests/test_csv_update.py @@ -24,7 +24,7 @@ "project1": { "artifact_layout": "project|run_index|run_id", "files": [".1.fastq.gz", ".2.fastq.gz", ".csv"], - "sites": ["subsite1.site1.project1", "site2.project1"], + "sites": ["subsite1.site1.project1", "site2.project1", "site1"], "bucket_policies": { "site_ingest": ["get", "put", "list", "delete"], "site_read": ["get", "list"], @@ -36,7 +36,11 @@ "ingest": { "name_layout": "{project}-{site}-{platform}-{test_flag}", "policy": "site_ingest", - } + }, + "results": { + "name_layout": "{project}-{site}-results", + "policy": "site_ingest", + }, }, "project_buckets": { "fake_files": { @@ -81,7 +85,7 @@ "project2": { "artifact_layout": "project|run_index|run_id", "files": [".1.fastq.gz", ".2.fastq.gz", ".csv"], - "sites": ["subsite1.site1.project2", "site2.project2"], + "sites": ["subsite1.site1.project2", "site2.project2", "site1"], "bucket_policies": { "site_ingest": ["get", "put", "list", "delete"], "site_read": ["get", "list"], @@ -92,7 +96,11 @@ "ingest": { "name_layout": "{project}-{site}-{platform}-{test_flag}", "policy": "site_ingest", - } + }, + "results": { + "name_layout": "{project}-{site}-results", + "policy": "site_ingest", + }, }, "csv_updates": True, "project_buckets": { diff --git a/tests/test_s3_matcher.py b/tests/test_s3_matcher.py index e3f31ed..b7e7ba9 100644 --- a/tests/test_s3_matcher.py +++ b/tests/test_s3_matcher.py @@ -19,7 +19,7 @@ "project1": { "artifact_layout": "project|run_index|run_id", "files": [".1.fastq.gz", ".2.fastq.gz", ".csv"], - "sites": ["subsite1.site1.project1", "site2.project1"], + "sites": ["subsite1.site1.project1", "site2.project1", "site1"], "bucket_policies": { "site_ingest": ["get", "put", "list", "delete"], "site_read": ["get", "list"], @@ -75,7 +75,7 @@ "project2": { "artifact_layout": "project|run_index|run_id", "files": [".1.fastq.gz", ".2.fastq.gz", ".csv"], - "sites": ["subsite1.site1.project2", "site2.project2"], + "sites": ["subsite1.site1.project2", "site2.project2", "site1"], "bucket_policies": { "site_ingest": ["get", "put", "list", "delete"], "site_read": ["get", "list"], From 38b5df7d00ab5403699c8ad16d1ecb540c0df200 Mon Sep 17 00:00:00 2001 From: Biowilko Date: Mon, 24 Aug 2026 12:38:22 +0100 Subject: [PATCH 2/3] fix test msg submission --- tests/test_integration.py | 20 ++++++++++---------- 1 file changed, 10 insertions(+), 10 deletions(-) diff --git a/tests/test_integration.py b/tests/test_integration.py index 428f806..6a79f1b 100644 --- a/tests/test_integration.py +++ b/tests/test_integration.py @@ -77,9 +77,9 @@ "s3SchemaVersion": "1.0", "configurationId": "inbound.s3", "bucket": { - "name": "mscape-birm-ont-prod", + "name": "mscape-birm.mscape-ont-prod", "ownerIdentity": {"principalId": "testuser"}, - "arn": "arn:aws:s3:::mscape-birm-ont-prod", + "arn": "arn:aws:s3:::mscape-birm.mscape-ont-prod", "id": "testdata", }, "object": { @@ -119,9 +119,9 @@ "s3SchemaVersion": "1.0", "configurationId": "inbound.s3", "bucket": { - "name": "mscape-birm-ont-prod", + "name": "mscape-birm.mscape-ont-prod", "ownerIdentity": {"principalId": "testuser"}, - "arn": "arn:aws:s3:::mscape-birm-ont-prod", + "arn": "arn:aws:s3:::mscape-birm.mscape-ont-prod", "id": "testdata", }, "object": { @@ -161,9 +161,9 @@ "s3SchemaVersion": "1.0", "configurationId": "inbound.s3", "bucket": { - "name": "mscape-birm-ont-prod", + "name": "mscape-birm.mscape-ont-prod", "ownerIdentity": {"principalId": "testuser"}, - "arn": "arn:aws:s3:::mscape-birm-ont-prod", + "arn": "arn:aws:s3:::mscape-birm.mscape-ont-prod", "id": "testdata", }, "object": { @@ -203,9 +203,9 @@ "s3SchemaVersion": "1.0", "configurationId": "inbound.s3", "bucket": { - "name": "mscape-birm-ont-prod", + "name": "mscape-birm.mscape-ont-prod", "ownerIdentity": {"principalId": "testuser"}, - "arn": "arn:aws:s3:::mscape-birm-ont-prod", + "arn": "arn:aws:s3:::mscape-birm.mscape-ont-prod", "id": "testdata", }, "object": { @@ -245,9 +245,9 @@ "s3SchemaVersion": "1.0", "configurationId": "inbound.s3", "bucket": { - "name": "mscape-birm-ont-prod", + "name": "mscape-birm.mscape-ont-prod", "ownerIdentity": {"principalId": "testuser"}, - "arn": "arn:aws:s3:::mscape-birm-ont-prod", + "arn": "arn:aws:s3:::mscape-birm.mscape-ont-prod", "id": "testdata", }, "object": { From c0c734d26b677a41059fac1a2071456fe801270b Mon Sep 17 00:00:00 2001 From: Biowilko Date: Mon, 24 Aug 2026 14:01:43 +0100 Subject: [PATCH 3/3] strip real codes from codebase --- docs/message-queue-payloads.md | 2 +- roz_scripts/utils/config.py | 2 +- tests/fixtures/test_config_adversarial.json | 4 ++-- tests/test_config.py | 6 +++--- 4 files changed, 7 insertions(+), 7 deletions(-) diff --git a/docs/message-queue-payloads.md b/docs/message-queue-payloads.md index 56c8412..bebce63 100644 --- a/docs/message-queue-payloads.md +++ b/docs/message-queue-payloads.md @@ -33,7 +33,7 @@ Emitted once all required files for an artifact have been seen. | Field | Type | Description | |-------|------|-------------| | `uuid` | string (UUID4) | Unique identifier for this match event | -| `site` | string | Submitting site name (short form, e.g. `"bham"`) | +| `site` | string | Submitting site name (short form, e.g. `"birm"`) | | `raw_site` | string | Submitting site as parsed from the bucket name (may include domain prefix) | | `uploaders` | array[string] | Deduplicated list of uploader IDs that contributed files | | `match_timestamp` | integer | Unix timestamp in nanoseconds when the match was made | diff --git a/roz_scripts/utils/config.py b/roz_scripts/utils/config.py index 6214771..6eb7c96 100644 --- a/roz_scripts/utils/config.py +++ b/roz_scripts/utils/config.py @@ -178,7 +178,7 @@ class BucketMatch: def short_site(site: str) -> str: """Derive the short form of a site name used for exchange names and dedup keys - e.g. "gpha.ukhsa.mscape" -> "ukhsa". Note this is a lossy heuristic that + e.g. "clinic1.trust1.mscape" -> "trust1". Note this is a lossy heuristic that can collide between distinct sites sharing a second-to-last dotted component - it exists here, centralised, purely to match pre-existing behaviour rather than as an endorsement of the scheme. diff --git a/tests/fixtures/test_config_adversarial.json b/tests/fixtures/test_config_adversarial.json index 0b1dc14..fbde6d4 100644 --- a/tests/fixtures/test_config_adversarial.json +++ b/tests/fixtures/test_config_adversarial.json @@ -6,8 +6,8 @@ "artifact_layout": "project|run_index|run_id", "sites": { "adversarialscape": "analysis", - "gpha.ukhsa.adversarialscape": "analysis", - "tarzet.ukhsa.adversarialscape": "uploader", + "clinic1.trust1.adversarialscape": "analysis", + "clinic2.trust1.adversarialscape": "uploader", "hyphen-site.adversarialscape": "uploader", "collide.adversarialscape": "uploader", "collide": "uploader" diff --git a/tests/test_config.py b/tests/test_config.py index a1a6ac5..a8197b3 100644 --- a/tests/test_config.py +++ b/tests/test_config.py @@ -21,7 +21,7 @@ class test_short_site(unittest.TestCase): def test_dotted_site_returns_second_to_last_component(self): - self.assertEqual(short_site("gpha.ukhsa.mscape"), "ukhsa") + self.assertEqual(short_site("clinic1.trust1.mscape"), "trust1") def test_bare_site_returns_itself(self): self.assertEqual(short_site("mscape"), "mscape") @@ -134,9 +134,9 @@ def test_hyphenated_site_parses_correctly(self): self.assertEqual(match.test_flag, "prod") def test_multi_level_dotted_site(self): - bucket_name = "adversarialscape-gpha.ukhsa.adversarialscape-illumina-prod" + bucket_name = "adversarialscape-clinic1.trust1.adversarialscape-illumina-prod" match = parse_bucket_name(self.config, bucket_name) - self.assertEqual(match.site, "gpha.ukhsa.adversarialscape") + self.assertEqual(match.site, "clinic1.trust1.adversarialscape") def test_illumina_se_platform_not_confused_with_illumina(self): bucket_name = "adversarialscape-adversarialscape-illumina.se-prod"