Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
40 changes: 39 additions & 1 deletion digital_land/phase/lookup.py
Original file line number Diff line number Diff line change
Expand Up @@ -33,13 +33,15 @@ def __init__(
issue_log=None,
operational_issue_log=None,
entity_range=[],
lookup_rules=None,
):
self.lookups = lookups
self.redirect_lookups = redirect_lookups
self.issues = issue_log
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 = {}
Expand All @@ -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"}
Expand Down Expand Up @@ -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:
Expand Down Expand Up @@ -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 {}
Expand Down Expand Up @@ -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
Expand Down
62 changes: 62 additions & 0 deletions digital_land/pipeline/main.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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()

Expand Down Expand Up @@ -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", "")
Expand Down Expand Up @@ -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

Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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,
),
]
)
Expand Down Expand Up @@ -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(
Expand Down
15 changes: 15 additions & 0 deletions digital_land/schema.py
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand Down
94 changes: 94 additions & 0 deletions tests/unit/phase/test_lookup.py
Original file line number Diff line number Diff line change
Expand Up @@ -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):
Expand Down
Loading