From 993ed6d3e31d8b41d6b6d69c6c3f8ff38b08fa14 Mon Sep 17 00:00:00 2001 From: Tom Brooks Date: Tue, 30 Jun 2026 16:43:27 +0100 Subject: [PATCH] feat: add new lookup-rule for entity numbers --- digital_land/phase/lookup.py | 40 ++++++++- digital_land/pipeline/main.py | 62 ++++++++++++++ digital_land/schema.py | 15 ++++ tests/unit/phase/test_lookup.py | 94 +++++++++++++++++++++ tests/unit/test_pipeline.py | 144 ++++++++++++++++++++++++++++++++ 5 files changed, 354 insertions(+), 1 deletion(-) diff --git a/digital_land/phase/lookup.py b/digital_land/phase/lookup.py index 758245c1d..c5dfb712d 100644 --- a/digital_land/phase/lookup.py +++ b/digital_land/phase/lookup.py @@ -33,6 +33,7 @@ def __init__( issue_log=None, operational_issue_log=None, entity_range=[], + lookup_rules=None, ): self.lookups = lookups self.redirect_lookups = redirect_lookups @@ -40,6 +41,7 @@ def __init__( self.operational_issues = operational_issue_log self.reverse_lookups = self.build_reverse_lookups() self.entity_range = entity_range + self.lookup_rules = lookup_rules or [] def build_reverse_lookups(self): reverse_lookups = {} @@ -52,6 +54,27 @@ def build_reverse_lookups(self): def lookup(self, **kwargs): return self.lookups.get(key(**kwargs), "") + def lookup_rule(self, prefix="", organisation="", reference=""): + try: + ref_int = int(reference) + except (ValueError, TypeError): + return "" + + prefix_norm = normalise(prefix) + org_norm = normalise(organisation) + + for rule in self.lookup_rules: + if normalise(rule["prefix"]) != prefix_norm: + continue + rule_org = normalise(rule.get("organisation", "")) + if rule_org and rule_org != org_norm: + continue + candidate = ref_int + rule["offset"] + if rule["entity-minimum"] <= candidate <= rule["entity-maximum"]: + return str(candidate) + + return "" + def check_associated_organisation(self, entity): if entity in self.reverse_lookups: keywords = {"authority", "development", "government"} @@ -85,6 +108,11 @@ def get_entity(self, block): reference=reference, ) or self.lookup(prefix=prefix, reference=reference) + or self.lookup_rule( + prefix=prefix, + organisation=organisation, + reference=reference, + ) ) if entity and self.entity_range: @@ -178,8 +206,11 @@ def __init__( issue_log=None, odp_collections=None, package_prefixes=None, + lookup_rules=None, ): - super().__init__(lookups, redirect_lookups, issue_log) + super().__init__( + lookups, redirect_lookups, issue_log, lookup_rules=lookup_rules + ) self.entity_field = "reference-entity" self.odp_collections = odp_collections self.package_prefixes = package_prefixes or {} @@ -232,6 +263,13 @@ def process(self, stream): ) find_entity = self.check_associated_organisation(find_entity) + if not find_entity: + find_entity = self.lookup_rule( + prefix=prefix, + organisation=organisation, + reference=reference, + ) + if not find_entity or ( str(find_entity) in self.redirect_lookups and int(self.redirect_lookups[str(find_entity)].get("status", 0)) == 410 diff --git a/digital_land/pipeline/main.py b/digital_land/pipeline/main.py index d320b9b2a..c259ce248 100644 --- a/digital_land/pipeline/main.py +++ b/digital_land/pipeline/main.py @@ -109,6 +109,7 @@ def __init__(self, path, dataset, specification=None, config=None): self.concat = {} self.migrate = {} self.lookup = {} + self.lookup_rule = {} self.redirect_lookup = {} self.specification = specification @@ -124,6 +125,7 @@ def __init__(self, path, dataset, specification=None, config=None): self.load_combine_fields() self.load_migrate() self.load_lookup() + self.load_lookup_rule() self.load_redirect_lookup() self.load_filter() @@ -292,6 +294,56 @@ def load_lookup(self): ) ] = row["entity"] + def load_lookup_rule(self): + # optional: collections without a lookup-rule.csv simply have no rules + rows = list(self.file_reader("lookup-rule.csv")) + if not rows: + return + + # validate columns against the schema; "pipeline" is an accepted + # legacy alias for the prefix/dataset column + allowed_fields = set(Schema("lookup-rule").fieldnames) | {"pipeline"} + observed_fields = set().union(*(row.keys() for row in rows)) + unexpected = observed_fields - allowed_fields + if unexpected: + raise RuntimeError( + "unexpected columns in lookup-rule.csv: " + f"{', '.join(sorted(unexpected))}" + ) + + rule_count = 0 + for row in rows: + prefix = ( + row.get("prefix", "") + or row.get("dataset", "") + or row.get("pipeline", "") + ) + organisation = row.get("organisation", "").replace( + "local-authority-eng", "local-authority" + ) + resource = row.get("resource", "") + + if ( + row.get("offset", "") + and row.get("entity-minimum", "") + and row.get("entity-maximum", "") + ): + rule = { + "prefix": prefix, + "organisation": organisation, + "offset": int(row["offset"]), + "entity-minimum": int(row["entity-minimum"]), + "entity-maximum": int(row["entity-maximum"]), + } + self.lookup_rule.setdefault(resource, []).append(rule) + rule_count += 1 + + if rule_count == 0: + logging.warning( + "lookup-rule.csv was found but produced no valid rules; " + "check the offset, entity-minimum and entity-maximum columns" + ) + def load_redirect_lookup(self): for row in self.file_reader("old-entity.csv"): old_entity = row.get("old-entity", "") @@ -423,6 +475,12 @@ def lookups(self, resource=None): d.update(self.lookup.get(resource, {})) return d + def lookup_rules(self, resource=None): + rules = list(self.lookup_rule.get("", [])) + if resource: + rules.extend(self.lookup_rule.get(resource, [])) + return rules + def redirect_lookups(self): return self.redirect_lookup @@ -577,6 +635,8 @@ def transform( concats = self.concatenations(resource, endpoints=endpoints) patches = self.patches(resource=resource, endpoints=endpoints) lookups = self.lookups(resource=resource) + lookup_rules = self.lookup_rules(resource=resource) + default_fields = self.default_fields(resource=resource, endpoints=endpoints) default_values = self.default_values(endpoints=endpoints) combine_fields = self.combine_fields(endpoints=endpoints) @@ -652,6 +712,7 @@ def transform( issue_log=self.issue_log, operational_issue_log=self.operational_issue_log, entity_range=[entity_range_min, entity_range_max], + lookup_rules=lookup_rules, ), ] ) @@ -691,6 +752,7 @@ def transform( issue_log=self.issue_log, odp_collections=self.specification.get_odp_collections(), package_prefixes=self.specification.get_package_prefixes(), + lookup_rules=lookup_rules, ), FactPrunePhase(), SavePhase( diff --git a/digital_land/schema.py b/digital_land/schema.py index 50add8238..fd60c1356 100644 --- a/digital_land/schema.py +++ b/digital_land/schema.py @@ -81,6 +81,21 @@ "end-date", ], }, + "lookup-rule": { + "key": "lookup-rule", + "fields": [ + "prefix", + "dataset", + "organisation", + "resource", + "offset", + "entity-minimum", + "entity-maximum", + "entry-date", + "start-date", + "end-date", + ], + }, "operational-issue": { "fields": [ "dataset", diff --git a/tests/unit/phase/test_lookup.py b/tests/unit/phase/test_lookup.py index bb38a0027..7fe963fab 100644 --- a/tests/unit/phase/test_lookup.py +++ b/tests/unit/phase/test_lookup.py @@ -147,6 +147,100 @@ def test_process_empty_prefix(self, get_lookup): assert output[0]["row"]["entity"] == "10" + def test_lookup_rule_valid_match(self): + lookup_rules = [ + { + "prefix": "dataset", + "organisation": "", + "offset": 1000, + "entity-minimum": 1000, + "entity-maximum": 4000, + } + ] + phase = LookupPhase(lookup_rules=lookup_rules) + assert phase.lookup_rule(prefix="dataset", reference="100") == "1100" + + def test_lookup_rule_out_of_range(self): + lookup_rules = [ + { + "prefix": "dataset", + "organisation": "", + "offset": 1000, + "entity-minimum": 2000, + "entity-maximum": 4000, + } + ] + phase = LookupPhase(lookup_rules=lookup_rules) + assert phase.lookup_rule(prefix="dataset", reference="100") == "" + + def test_lookup_rule_non_integer_reference(self): + lookup_rules = [ + { + "prefix": "dataset", + "organisation": "", + "offset": 1000, + "entity-minimum": 1000, + "entity-maximum": 4000, + } + ] + phase = LookupPhase(lookup_rules=lookup_rules) + assert phase.lookup_rule(prefix="dataset", reference="abc") == "" + + def test_lookup_rule_no_matching_prefix(self): + lookup_rules = [ + { + "prefix": "other-dataset", + "organisation": "", + "offset": 1000, + "entity-minimum": 1000, + "entity-maximum": 4000, + } + ] + phase = LookupPhase(lookup_rules=lookup_rules) + assert phase.lookup_rule(prefix="dataset", reference="100") == "" + + def test_lookup_rule_with_organisation(self): + lookup_rules = [ + { + "prefix": "dataset", + "organisation": "local-authority:ABC", + "offset": 1000, + "entity-minimum": 1000, + "entity-maximum": 4000, + } + ] + phase = LookupPhase(lookup_rules=lookup_rules) + assert ( + phase.lookup_rule( + prefix="dataset", organisation="local-authority:ABC", reference="100" + ) + == "1100" + ) + assert ( + phase.lookup_rule( + prefix="dataset", organisation="local-authority:XYZ", reference="100" + ) + == "" + ) + + def test_lookup_rule_fallback_in_get_entity(self): + lookup_rules = [ + { + "prefix": "dataset", + "organisation": "", + "offset": 1000, + "entity-minimum": 1000, + "entity-maximum": 4000, + } + ] + phase = LookupPhase(lookup_rules=lookup_rules) + phase.entity_field = "entity" + block = { + "row": {"prefix": "dataset", "reference": "100", "organisation": ""}, + "entry-number": 1, + } + assert phase.get_entity(block) == "1100" + class TestPrintLookupPhase: def test_process_does_not_produce_new_lookup(self, get_input_stream, get_lookup): diff --git a/tests/unit/test_pipeline.py b/tests/unit/test_pipeline.py index 9d57c6abb..737896243 100755 --- a/tests/unit/test_pipeline.py +++ b/tests/unit/test_pipeline.py @@ -1,5 +1,6 @@ #!/usr/bin/env -S py.test -svv import io +import logging import pytest from digital_land.pipeline import EntityNumGen, Pipeline from digital_land.pipeline import Lookups @@ -367,6 +368,149 @@ def open_mock(file, *args, **kwargs): if new_lookup[0]["entity"] is None: assert True + def test_load_lookup_rule_with_all_fields(self, mocker): + def mock_file_reader(self, filepath): + if filepath == "lookup-rule.csv": + return [ + { + "prefix": "listed-building", + "organisation": "", + "offset": "1000000", + "entity-minimum": "1000000", + "entity-maximum": "2000000", + "resource": "", + } + ] + return [] + + mocker.patch("digital_land.pipeline.Pipeline.file_reader", mock_file_reader) + p = Pipeline("anything", "stuff") + rules = p.lookup_rules() + assert len(rules) == 1 + assert rules[0]["offset"] == 1000000 + assert rules[0]["entity-minimum"] == 1000000 + assert rules[0]["entity-maximum"] == 2000000 + assert rules[0]["prefix"] == "listed-building" + + def test_load_lookup_rule_missing_entity_minimum(self, mocker): + def mock_file_reader(self, filepath): + if filepath == "lookup-rule.csv": + return [ + { + "prefix": "listed-building", + "organisation": "", + "offset": "1000000", + "entity-minimum": "", + "entity-maximum": "2000000", + "resource": "", + } + ] + return [] + + mocker.patch("digital_land.pipeline.Pipeline.file_reader", mock_file_reader) + p = Pipeline("anything", "stuff") + assert len(p.lookup_rules()) == 0 + + def test_load_lookup_rule_missing_entity_maximum(self, mocker): + def mock_file_reader(self, filepath): + if filepath == "lookup-rule.csv": + return [ + { + "prefix": "listed-building", + "organisation": "", + "offset": "1000000", + "entity-minimum": "1000000", + "entity-maximum": "", + "resource": "", + } + ] + return [] + + mocker.patch("digital_land.pipeline.Pipeline.file_reader", mock_file_reader) + p = Pipeline("anything", "stuff") + assert len(p.lookup_rules()) == 0 + + def test_load_lookup_rule_unexpected_column(self, mocker): + def mock_file_reader(self, filepath): + if filepath == "lookup-rule.csv": + return [ + { + "prefix": "listed-building", + "offset": "1000000", + "entity-min": "1000000", # typo for entity-minimum + "entity-maximum": "2000000", + "resource": "", + } + ] + return [] + + mocker.patch("digital_land.pipeline.Pipeline.file_reader", mock_file_reader) + with pytest.raises(RuntimeError, match="unexpected columns"): + Pipeline("anything", "stuff") + + def test_load_lookup_rule_warns_when_no_valid_rules(self, mocker, caplog): + def mock_file_reader(self, filepath): + if filepath == "lookup-rule.csv": + return [ + { + "prefix": "listed-building", + "organisation": "", + "offset": "1000000", + "entity-minimum": "", # missing -> no rule built + "entity-maximum": "2000000", + "resource": "", + } + ] + return [] + + mocker.patch("digital_land.pipeline.Pipeline.file_reader", mock_file_reader) + with caplog.at_level(logging.WARNING): + p = Pipeline("anything", "stuff") + + assert len(p.lookup_rules()) == 0 + assert any("lookup-rule.csv" in r.message for r in caplog.records) + + def test_lookup_rules_resource_scoping(self, mocker): + def mock_file_reader(self, filepath): + if filepath == "lookup-rule.csv": + return [ + { + "prefix": "listed-building", + "organisation": "", + "offset": "1000000", + "entity-minimum": "1000000", + "entity-maximum": "2000000", + "resource": "", # global rule + }, + { + "prefix": "listed-building", + "organisation": "", + "offset": "5000000", + "entity-minimum": "5000000", + "entity-maximum": "6000000", + "resource": "res1", # resource-specific rule + }, + ] + return [] + + mocker.patch("digital_land.pipeline.Pipeline.file_reader", mock_file_reader) + p = Pipeline("anything", "stuff") + + # no resource -> only the global rule + global_rules = p.lookup_rules() + assert len(global_rules) == 1 + assert global_rules[0]["offset"] == 1000000 + + # matching resource -> global rule plus the resource-specific one + scoped_rules = p.lookup_rules(resource="res1") + assert len(scoped_rules) == 2 + assert {r["offset"] for r in scoped_rules} == {1000000, 5000000} + + # unrelated resource -> only the global rule + other_rules = p.lookup_rules(resource="other") + assert len(other_rules) == 1 + assert other_rules[0]["offset"] == 1000000 + @pytest.fixture def pipeline(self, mocker): def mock_file_reader(self, filename):