368 lines
15 KiB
Python
368 lines
15 KiB
Python
|
|
#!/usr/bin/env python3
|
||
|
|
"""
|
||
|
|
Migrate the Pernod Ricard Data MetaModel and the SODH instances to v2.2.
|
||
|
|
|
||
|
|
USAGE
|
||
|
|
python3 migrate_tbox_v2_2.py --ontology ... --instances ...
|
||
|
|
python3 migrate_tbox_v2_2.py --ontology ... --instances ... --apply
|
||
|
|
|
||
|
|
v2.2 turns the actor nodes into governance ASSIGNMENTS: a stable post carrying
|
||
|
|
an identifier, a domain and a name, plus the person holding it as a value. It
|
||
|
|
completes the twelve that existed, creates the ten sub-domain owners, the
|
||
|
|
product owner, and one steward assignment per business object.
|
||
|
|
|
||
|
|
DESIGN
|
||
|
|
EV-004 dry run is the default
|
||
|
|
EV-005 every edit goes through rdflib
|
||
|
|
EV-006 guards test the target state, so a replay reports zero change
|
||
|
|
EV-015 an execution log is written for every attempt
|
||
|
|
"""
|
||
|
|
import argparse
|
||
|
|
import datetime
|
||
|
|
import hashlib
|
||
|
|
import json
|
||
|
|
import os
|
||
|
|
import re
|
||
|
|
import sys
|
||
|
|
from collections import OrderedDict
|
||
|
|
|
||
|
|
import yaml
|
||
|
|
from rdflib import Graph, Literal, Namespace, RDF, RDFS, OWL, URIRef
|
||
|
|
|
||
|
|
HERE = os.path.dirname(os.path.abspath(__file__))
|
||
|
|
SPEC = os.path.join(HERE, "assignments_v2_2.yaml")
|
||
|
|
|
||
|
|
|
||
|
|
class Report(object):
|
||
|
|
def __init__(self):
|
||
|
|
self.steps = OrderedDict()
|
||
|
|
self.notes = []
|
||
|
|
|
||
|
|
def add(self, step, n=1, detail=None):
|
||
|
|
self.steps[step] = self.steps.get(step, 0) + n
|
||
|
|
if detail:
|
||
|
|
self.notes.append("%s: %s" % (step, detail))
|
||
|
|
|
||
|
|
@property
|
||
|
|
def total(self):
|
||
|
|
return sum(self.steps.values())
|
||
|
|
|
||
|
|
|
||
|
|
def derive_label(name):
|
||
|
|
"""TN-018 and TN-023: a label is derived from the local name, never typed."""
|
||
|
|
spaced = re.sub(r"(?<=[a-z0-9])(?=[A-Z])|(?<=[A-Z])(?=[A-Z][a-z])", " ", name)
|
||
|
|
return spaced[0].lower() + spaced[1:]
|
||
|
|
|
||
|
|
|
||
|
|
def md5(path):
|
||
|
|
h = hashlib.md5()
|
||
|
|
with open(path, "rb") as fh:
|
||
|
|
for chunk in iter(lambda: fh.read(65536), b""):
|
||
|
|
h.update(chunk)
|
||
|
|
return h.hexdigest()
|
||
|
|
|
||
|
|
|
||
|
|
def rename_node(g, old, new, report, step):
|
||
|
|
"""EV-001 applied to an instance: rewrite every triple naming the node."""
|
||
|
|
if (new, None, None) in g and (old, None, None) not in g:
|
||
|
|
return 0
|
||
|
|
n = 0
|
||
|
|
for s, p, o in list(g):
|
||
|
|
ns, no = (new if s == old else s), (new if o == old else o)
|
||
|
|
if (ns, no) != (s, o):
|
||
|
|
g.remove((s, p, o))
|
||
|
|
g.add((ns, p, no))
|
||
|
|
n += 1
|
||
|
|
report.add(step, n)
|
||
|
|
return n
|
||
|
|
|
||
|
|
|
||
|
|
def assign(g, pr, node, role, identifier, name, held_by, domain, spec, report, step):
|
||
|
|
"""Write one assignment. The guard tests the identifier, so a replay is a no-op."""
|
||
|
|
if (node, pr.identifier, Literal(identifier)) in g:
|
||
|
|
return False
|
||
|
|
g.add((node, RDF.type, role))
|
||
|
|
g.remove((node, pr.identifier, None))
|
||
|
|
g.add((node, pr.identifier, Literal(identifier)))
|
||
|
|
g.remove((node, pr.canonicalName, None))
|
||
|
|
g.add((node, pr.canonicalName, Literal(name)))
|
||
|
|
g.remove((node, pr.personName, None))
|
||
|
|
if held_by:
|
||
|
|
g.add((node, pr.personName, Literal(held_by)))
|
||
|
|
g.remove((node, pr.status, None))
|
||
|
|
g.add((node, pr.status, Literal(spec["status"])))
|
||
|
|
g.remove((node, pr.version, None))
|
||
|
|
g.add((node, pr.version, Literal(spec["version"])))
|
||
|
|
if domain is not None:
|
||
|
|
g.remove((node, pr.ownedByDomain, None))
|
||
|
|
g.add((node, pr.ownedByDomain, domain))
|
||
|
|
report.add(step, 1)
|
||
|
|
return True
|
||
|
|
|
||
|
|
|
||
|
|
def check_version(g, spec, path):
|
||
|
|
"""EV-010 applied to a migration script: refuse to run against the wrong input.
|
||
|
|
|
||
|
|
A script announcing 2.1 -> 2.2 will happily work on a file still in 2.0 and
|
||
|
|
report success, leaving a state that is neither version. The same silent
|
||
|
|
failure the shapes are guarded against, one layer down.
|
||
|
|
"""
|
||
|
|
ns = spec["meta"]["namespace"]
|
||
|
|
expected = URIRef(ns.rstrip("/") + "/" + spec["meta"]["from_version"])
|
||
|
|
target = URIRef(ns.rstrip("/") + "/" + spec["meta"]["to_version"])
|
||
|
|
found = set(g.objects(None, OWL.versionIRI))
|
||
|
|
if target in found:
|
||
|
|
print("ALREADY AT %s — nothing to do." % spec["meta"]["to_version"])
|
||
|
|
sys.exit(0)
|
||
|
|
if expected not in found:
|
||
|
|
print("ABORT — %s declares %s; this script migrates from %s."
|
||
|
|
% (os.path.basename(path),
|
||
|
|
", ".join(sorted(str(f) for f in found)) or "no version",
|
||
|
|
expected))
|
||
|
|
sys.exit(1)
|
||
|
|
|
||
|
|
|
||
|
|
def migrate(onto, data, spec, report):
|
||
|
|
ns, ins = spec["meta"]["namespace"], spec["meta"]["instance_namespace"]
|
||
|
|
pr, ex = Namespace(ns), Namespace(ins)
|
||
|
|
|
||
|
|
# 1 — declare heldBy on Actor (procedure C1)
|
||
|
|
for d in spec.get("declare_properties") or []:
|
||
|
|
term = pr[d["term"]]
|
||
|
|
if (term, RDF.type, getattr(OWL, d["kind"])) not in onto:
|
||
|
|
onto.add((term, RDF.type, getattr(OWL, d["kind"])))
|
||
|
|
onto.add((term, RDFS.label, Literal(derive_label(d["term"]))))
|
||
|
|
onto.add((term, RDFS.domain, pr[d["domain"]]))
|
||
|
|
onto.add((term, RDFS.range, URIRef(d["range"])))
|
||
|
|
onto.add((term, RDFS.comment, Literal(" ".join(d["comment"].split()))))
|
||
|
|
report.add("properties_declared", 1, d["term"])
|
||
|
|
|
||
|
|
# 2 — the short form of Data Steward is DST, and the rulebook now says so
|
||
|
|
for u in spec.get("display_updates") or []:
|
||
|
|
term = pr[u["term"]]
|
||
|
|
if (term, pr.acronym, Literal(u["acronym"])) not in onto:
|
||
|
|
onto.remove((term, pr.acronym, None))
|
||
|
|
onto.add((term, pr.acronym, Literal(u["acronym"])))
|
||
|
|
report.add("display_updated", 1, "%s -> %s" % (u["term"], u["acronym"]))
|
||
|
|
|
||
|
|
# 3 — rename the instance nodes that carried the wrong short form
|
||
|
|
for r in spec.get("instance_renames") or []:
|
||
|
|
rename_node(data, ex[r["from"]], ex[r["to"]], report, "instances_renamed")
|
||
|
|
|
||
|
|
# 4 — complete the assignments that already existed
|
||
|
|
for c in spec.get("complete") or []:
|
||
|
|
node = ex[c["node"]]
|
||
|
|
if (node, None, None) not in data:
|
||
|
|
report.add("complete_missing", 1, c["node"])
|
||
|
|
continue
|
||
|
|
role = next(iter(data.objects(node, RDF.type)), None)
|
||
|
|
assign(data, pr, node, role, c["identifier"], c["name"],
|
||
|
|
c["held_by"], ex[c["domain"]], spec, report, "assignments_completed")
|
||
|
|
|
||
|
|
# 5 — sub-domain owners
|
||
|
|
sdo_of = {}
|
||
|
|
for s in spec.get("create_sub_domain_owners") or []:
|
||
|
|
node = ex[s["node"]]
|
||
|
|
sdo_of[ex[s["sub_domain"]]] = node
|
||
|
|
if assign(data, pr, node, pr.DataSubDomainOwner, s["identifier"],
|
||
|
|
spec["sub_domain_owner_name"], s["held_by"], ex[s["domain"]],
|
||
|
|
spec, report, "sub_domain_owners_created"):
|
||
|
|
data.add((ex[s["sub_domain"]], pr.hasSubDomainOwner, node))
|
||
|
|
|
||
|
|
# 6 — product owners
|
||
|
|
for p in spec.get("create_product_owners") or []:
|
||
|
|
node = ex[p["node"]]
|
||
|
|
product = ex[p["product"]]
|
||
|
|
if (product, None, None) not in data:
|
||
|
|
report.add("product_missing", 1, p["product"])
|
||
|
|
continue
|
||
|
|
if assign(data, pr, node, pr.DataProductOwner, p["identifier"],
|
||
|
|
spec["product_owner_name"], p["held_by"], ex[p["domain"]],
|
||
|
|
spec, report, "product_owners_created"):
|
||
|
|
data.add((product, pr.hasProductOwner, node))
|
||
|
|
|
||
|
|
# 7 — create the business objects that close a governance chain
|
||
|
|
for b in spec.get("create_business_objects") or []:
|
||
|
|
node = ex[b["node"]]
|
||
|
|
if (node, pr.identifier, Literal(b["identifier"])) in data:
|
||
|
|
continue
|
||
|
|
data.add((node, RDF.type, pr.BusinessObject))
|
||
|
|
data.add((node, pr.identifier, Literal(b["identifier"])))
|
||
|
|
data.add((node, pr.canonicalName, Literal(b["name"])))
|
||
|
|
data.add((node, pr.isAbout, ex[b["is_about"]]))
|
||
|
|
data.add((node, pr.belongsTo, ex[b["sub_domain"]]))
|
||
|
|
data.add((node, pr.ownedByDomain, ex[b["domain"]]))
|
||
|
|
data.add((node, pr.status, Literal(spec["status"])))
|
||
|
|
data.add((node, pr.version, Literal(spec["version"])))
|
||
|
|
report.add("business_objects_created", 1, b["node"])
|
||
|
|
|
||
|
|
# 8 — one steward assignment per business object, derived from the graph
|
||
|
|
holders = spec.get("steward_holders") or {}
|
||
|
|
steward_of = {}
|
||
|
|
for bo in sorted(data.subjects(RDF.type, pr.BusinessObject), key=str):
|
||
|
|
ident = next((str(v) for v in data.objects(bo, pr.identifier)), None)
|
||
|
|
if not ident:
|
||
|
|
report.add("steward_no_identifier", 1, str(bo).replace(ins, ""))
|
||
|
|
continue
|
||
|
|
suffix = ident.split("-", 1)[1] # 06.01-001
|
||
|
|
domain_key = "DD_" + suffix.split(".", 1)[0] # DD_06
|
||
|
|
holder = holders.get(domain_key)
|
||
|
|
if holder is None:
|
||
|
|
report.add("steward_no_holder", 1, domain_key)
|
||
|
|
continue
|
||
|
|
node = ex["DST_" + suffix.replace(".", "_").replace("-", "_")]
|
||
|
|
steward_of[bo] = node
|
||
|
|
if assign(data, pr, node, pr.DataSteward, "DST-" + suffix,
|
||
|
|
spec["steward_name"], holder, ex[domain_key],
|
||
|
|
spec, report, "stewards_created"):
|
||
|
|
for _, _, old in list(data.triples((bo, pr.monitoredBy, None))):
|
||
|
|
data.remove((bo, pr.monitoredBy, old))
|
||
|
|
data.add((bo, pr.monitoredBy, node))
|
||
|
|
|
||
|
|
# 9 — resolve every ownedBy pointing at a per-domain steward node
|
||
|
|
res = spec.get("owned_by_resolution") or {}
|
||
|
|
old_stewards = {ex[r["from"]] for r in (spec.get("instance_renames") or [])} | \
|
||
|
|
{ex[r["to"]] for r in (spec.get("instance_renames") or [])}
|
||
|
|
|
||
|
|
def business_object_of(subject):
|
||
|
|
"""Follow the declared chains until a business object is reached."""
|
||
|
|
for chain in res.get("chains") or []:
|
||
|
|
prop = pr[chain["via"]]
|
||
|
|
if chain["direction"] == "inbound":
|
||
|
|
for candidate in data.subjects(prop, subject):
|
||
|
|
if (candidate, RDF.type, pr.BusinessObject) in data:
|
||
|
|
return candidate
|
||
|
|
else:
|
||
|
|
for target in data.objects(subject, prop):
|
||
|
|
if (target, RDF.type, pr.BusinessObject) in data:
|
||
|
|
return target
|
||
|
|
for candidate in data.subjects(pr.isAbout, target):
|
||
|
|
if (candidate, RDF.type, pr.BusinessObject) in data:
|
||
|
|
return candidate
|
||
|
|
return None
|
||
|
|
|
||
|
|
remove_on = {pr[c] for c in (res.get("remove_on") or [])}
|
||
|
|
for subject, _, target in list(data.triples((None, pr.ownedBy, None))):
|
||
|
|
if target not in old_stewards:
|
||
|
|
continue
|
||
|
|
types = set(data.objects(subject, RDF.type))
|
||
|
|
if types & remove_on:
|
||
|
|
data.remove((subject, pr.ownedBy, target))
|
||
|
|
report.add("ownedby_removed", 1)
|
||
|
|
continue
|
||
|
|
if pr.DataSubDomain in types and res.get("repoint_sub_domains"):
|
||
|
|
owner = sdo_of.get(subject)
|
||
|
|
if owner is None:
|
||
|
|
report.add("ownedby_unresolved", 1,
|
||
|
|
"%s has no sub-domain owner" % str(subject).replace(ins, ""))
|
||
|
|
continue
|
||
|
|
data.remove((subject, pr.ownedBy, target))
|
||
|
|
data.add((subject, pr.ownedBy, owner))
|
||
|
|
report.add("ownedby_repointed_to_sdo", 1)
|
||
|
|
continue
|
||
|
|
bo = business_object_of(subject)
|
||
|
|
owner = steward_of.get(bo) if bo is not None else None
|
||
|
|
if owner is None:
|
||
|
|
report.add("ownedby_unresolved", 1,
|
||
|
|
"%s reaches no business object" % str(subject).replace(ins, ""))
|
||
|
|
continue
|
||
|
|
data.remove((subject, pr.ownedBy, target))
|
||
|
|
data.add((subject, pr.ownedBy, owner))
|
||
|
|
report.add("ownedby_repointed_to_steward", 1)
|
||
|
|
|
||
|
|
# 10 — withdraw the per-domain steward nodes, once nothing names them (EV-003)
|
||
|
|
if spec.get("retire_domain_stewards"):
|
||
|
|
for r in spec.get("instance_renames") or []:
|
||
|
|
node = ex[r["to"]]
|
||
|
|
referrers = {s for s in data.subjects(None, node)}
|
||
|
|
if referrers:
|
||
|
|
report.add("steward_still_referenced", 1,
|
||
|
|
"%s cited by %d subject(s)" % (r["to"], len(referrers)))
|
||
|
|
continue
|
||
|
|
n = 0
|
||
|
|
for t in list(data.triples((node, None, None))):
|
||
|
|
data.remove(t)
|
||
|
|
n += 1
|
||
|
|
if n:
|
||
|
|
report.add("domain_stewards_retired", 1, r["to"])
|
||
|
|
|
||
|
|
# 11 — bump the version
|
||
|
|
target = URIRef(ns.rstrip("/") + "/" + spec["meta"]["to_version"])
|
||
|
|
for o in set(onto.subjects(RDF.type, OWL.Ontology)):
|
||
|
|
if (o, OWL.versionIRI, target) not in onto:
|
||
|
|
onto.remove((o, OWL.versionIRI, None))
|
||
|
|
onto.add((o, OWL.versionIRI, target))
|
||
|
|
report.add("version_bumped", 1, str(target))
|
||
|
|
|
||
|
|
|
||
|
|
def main():
|
||
|
|
ap = argparse.ArgumentParser()
|
||
|
|
ap.add_argument("--apply", action="store_true")
|
||
|
|
ap.add_argument("--ontology", required=True)
|
||
|
|
ap.add_argument("--instances", required=True)
|
||
|
|
ap.add_argument("--log-dir", default=os.path.join(HERE, "logs"))
|
||
|
|
args = ap.parse_args()
|
||
|
|
|
||
|
|
for path in (args.ontology, args.instances):
|
||
|
|
if not os.path.exists(path):
|
||
|
|
print("MISSING INPUT: %s" % path)
|
||
|
|
sys.exit(1)
|
||
|
|
|
||
|
|
spec = yaml.safe_load(open(SPEC, encoding="utf-8"))
|
||
|
|
checksums = OrderedDict((p, md5(p)) for p in (args.ontology, args.instances))
|
||
|
|
|
||
|
|
onto, data = Graph(), Graph()
|
||
|
|
onto.parse(args.ontology, format="turtle")
|
||
|
|
data.parse(args.instances, format="turtle")
|
||
|
|
check_version(onto, spec, args.ontology)
|
||
|
|
|
||
|
|
before = (len(onto), len(data))
|
||
|
|
|
||
|
|
report = Report()
|
||
|
|
migrate(onto, data, spec, report)
|
||
|
|
after = (len(onto), len(data))
|
||
|
|
|
||
|
|
print("MIGRATION %s -> %s %s"
|
||
|
|
% (spec["meta"]["from_version"], spec["meta"]["to_version"],
|
||
|
|
"APPLY" if args.apply else "DRY RUN"))
|
||
|
|
print()
|
||
|
|
for step, n in report.steps.items():
|
||
|
|
print(" %-28s %6d" % (step, n))
|
||
|
|
print(" %-28s %6d" % ("total changes", report.total))
|
||
|
|
print()
|
||
|
|
for path, b, a in zip((args.ontology, args.instances), before, after):
|
||
|
|
print(" %-46s %6d -> %6d triples" % (os.path.basename(path), b, a))
|
||
|
|
if report.notes:
|
||
|
|
print()
|
||
|
|
for n in report.notes:
|
||
|
|
print(" note: " + n)
|
||
|
|
|
||
|
|
os.makedirs(args.log_dir, exist_ok=True)
|
||
|
|
stamp = datetime.datetime.now().strftime("%Y%m%dT%H%M%S")
|
||
|
|
log_path = os.path.join(args.log_dir, "migration_v2_2_%s.json" % stamp)
|
||
|
|
json.dump({"attempt_timestamp": datetime.datetime.now().isoformat(timespec="seconds"),
|
||
|
|
"mode": "apply" if args.apply else "dry-run",
|
||
|
|
"from_version": spec["meta"]["from_version"],
|
||
|
|
"to_version": spec["meta"]["to_version"],
|
||
|
|
"input_checksums": checksums,
|
||
|
|
"steps": report.steps, "total_changes": report.total,
|
||
|
|
"triples_before": dict(zip((args.ontology, args.instances), before)),
|
||
|
|
"triples_after": dict(zip((args.ontology, args.instances), after)),
|
||
|
|
"notes": report.notes},
|
||
|
|
open(log_path, "w", encoding="utf-8"), indent=2)
|
||
|
|
print("\n log: %s" % log_path)
|
||
|
|
|
||
|
|
if not args.apply:
|
||
|
|
print("\n DRY RUN — nothing written. Re-run with --apply when the counts "
|
||
|
|
"above are what you expect.")
|
||
|
|
return
|
||
|
|
|
||
|
|
onto.serialize(destination=args.ontology, format="turtle")
|
||
|
|
data.serialize(destination=args.instances, format="turtle")
|
||
|
|
print("\n written. Replay this script now: a second run must report zero "
|
||
|
|
"changes (EV-006).")
|
||
|
|
|
||
|
|
|
||
|
|
if __name__ == "__main__":
|
||
|
|
main()
|