import os import json import glob import re import csv from collections import Counter, defaultdict from pathlib import Path # Fallback acronym mappings used only when canonical_map.json is absent COMMON_ACRONYM_MAP = { "wdi": "World Development Indicators", "dhs": "Demographic and Health Survey", "lsms": "Living Standards Measurement Study", "pwt": "Penn World Table", "mics": "Multiple Indicator Cluster Survey", "wvs": "World Values Survey", "lfs": "Labour Force Survey", "hces": "Household Consumption and Expenditure Survey", "mxfls": "Mexican Family Life Survey", "psid": "Panel Study of Income Dynamics", "bhps": "British Household Panel Survey", "soep": "German Socio-Economic Panel", "hilda": "Household, Income and Labour Dynamics in Australia", "gsoep": "German Socio-Economic Panel Study", "ipums": "Integrated Public Use Microdata Series", } def load_canonical_map() -> dict: """ Load the harmonization-produced canonical_map.json if available. Keys are normalized variant strings; values are the formal canonical names (e.g. 'Demographic and Health Survey (DHS)'). Returns an empty dict if the file does not yet exist. """ canonical_map_path = Path(__file__).parent / "canonical_map.json" if canonical_map_path.exists(): with open(canonical_map_path, "r", encoding="utf-8") as f: mapping = json.load(f) print(f" Loaded harmonization canonical map: {len(mapping)} variant -> canonical entries.") return mapping else: print(" canonical_map.json not found – falling back to basic acronym lookup.") return {} def clean_name(name): """Normalize dataset name for string matching.""" if not name: return "" s = name.lower().strip() s = re.sub(r'^(the|a|an|our)\s+', '', s) s = re.sub(r'[^\w\s]', ' ', s) s = re.sub(r'\s+', ' ', s).strip() return s def clean_acronym(acronym): """Normalize acronym.""" if not acronym: return "" s = acronym.upper().strip() s = re.sub(r'[^\w]', '', s) return s def sanitize_author_id(name): """Generate a safe ID from an author's name.""" s = name.strip() s = re.sub(r'[^\w\s-]', '', s) s = re.sub(r'[\s]+', '_', s) return s.lower() def main(): base_dir = Path(__file__).parent.parent standardized_base = Path("/Users/rafaelmacalaba/WBG/fetch_prwp/data/standardized_outputs") output_json_path = base_dir / "data" / "dashboard_data.json" output_csv_path = base_dir / "data" / "deduplicated_datasets.csv" graph_base = base_dir / "data" / "graph_database" # Load the harmonization canonical map (produced by run_harmonization.py) print("Step 0: Loading harmonization canonical map...") harmonization_map = load_canonical_map() print("Step 1: Reading all standardized JSON files...") all_json_files = glob.glob(str(standardized_base / "batch_*" / "*.json")) print(f" Found {len(all_json_files)} standardized JSON files.") # Data structures to hold records papers = [] # Named mentions: formal, identifiable dataset names → entity resolution + graph named_mentions = [] # Descriptive mentions: category/type references → data practice breakdown descriptive_mentions = [] # Vague mentions ("the data", "survey data") vague_mentions = [] # Author structures all_authors = set() paper_to_authors = defaultdict(list) # Track mapping of acronym -> full names with frequencies (named only) acronym_to_names = defaultdict(Counter) for filepath in all_json_files: try: with open(filepath, "r", encoding="utf-8") as f: doc = json.load(f) except Exception as e: print(f" [Error] Failed to read {os.path.basename(filepath)}: {e}") continue metadata = doc.get("metadata") or {} model_extractions = doc.get("model_extractions") or [] # Parse paper metadata paper_id = metadata.get("id") title = metadata.get("display_title", "Untitled Document").strip() pdf_url = metadata.get("pdfurl", "").strip() # Resolve year from publication date doc_date = metadata.get("docdt") or metadata.get("last_modified_date") or "" year = "Unknown" match = re.search(r"\b(19\d{2}|20\d{2})\b", doc_date) if match: year = int(match.group(1)) # Extract authors authors_dict = metadata.get("authors") or {} authors_list = [] if isinstance(authors_dict, dict): for k, v in authors_dict.items(): if isinstance(v, dict) and "author" in v: authors_list.append(v["author"].strip()) elif isinstance(v, str): authors_list.append(v.strip()) elif isinstance(authors_dict, list): authors_list = [a.get("author").strip() if isinstance(a, dict) else a.strip() for a in authors_dict] for author in authors_list: if author: all_authors.add(author) paper_to_authors[paper_id].append(author) authors_str = ", ".join(authors_list) if authors_list else "Unknown" # Track all mentions in this paper (any specificity) for the has_data flag paper_mention_count = 0 for extraction in model_extractions: if extraction.get("classifier_skipped", False): continue page_num = extraction.get("page", 0) + 1 datasets = extraction.get("datasets") or [] for ds in datasets: mention = ds.get("mention_name", {}).get("text", "").strip() acronym = ds.get("acronym", {}).get("text", "").strip() producer = ds.get("producer", {}).get("text", "").strip() geography = ds.get("geography", {}).get("text", "").strip() typology = ds.get("typology_tag", {}).get("text", "").strip() specificity = (ds.get("specificity_tag") or {}).get("text", "").strip().lower() usage_context = ds.get("usage_context", {}).get("text", "").strip() is_used_val = ds.get("is_used", {}).get("text", "").strip() confidence = ds.get("mention_name", {}).get("confidence", 0.0) if not mention: continue paper_mention_count += 1 record = { "paper_id": paper_id, "paper_title": title, "paper_year": year, "page": page_num, "mention": mention, "acronym": acronym, "producer": producer, "geography": geography, "typology": typology, "specificity": specificity, "usage_context": usage_context, "is_used": is_used_val, "confidence": confidence, } if specificity == "named": named_mentions.append(record) # Build acronym co-occurrence map for fallback canonical resolution cleaned_ac = clean_acronym(acronym) if cleaned_ac and len(mention) > len(acronym) and len(mention) > 5: acronym_to_names[cleaned_ac][mention] += 1 elif specificity == "descriptive": descriptive_mentions.append(record) elif specificity == "vague": vague_mentions.append(record) papers.append({ "id": paper_id, "title": title, "year": year, "authors": authors_str, "pdf_url": pdf_url, # has_data = True if the paper has any named or descriptive mention (not vague) "has_data": paper_mention_count > 0, "mention_count": paper_mention_count, }) print(f" Named mentions: {len(named_mentions)} | Descriptive: {len(descriptive_mentions)} | Vague: {len(vague_mentions)}") print("Step 2: Resolving canonical names via harmonization map (named mentions only)...") # Build a fallback dynamic acronym map from co-occurrence patterns in corpus dynamic_acronym_map = {} for ac, name_counts in acronym_to_names.items(): best_name = name_counts.most_common(1)[0][0] dynamic_acronym_map[ac.lower()] = best_name # Merge static fallback entries for ac, name in COMMON_ACRONYM_MAP.items(): if ac.lower() not in dynamic_acronym_map: dynamic_acronym_map[ac.lower()] = name harmonization_hits = 0 fallback_hits = 0 # Map NAMED mentions to canonical names # Priority: (1) harmonization_map, (2) dynamic/static acronym map, (3) most-frequent raw string raw_to_canonical = {} cleaned_groups = defaultdict(list) for m in named_mentions: cleaned = clean_name(m["mention"]) cleaned_groups[cleaned].append(m) for cleaned, group in cleaned_groups.items(): raw_names = [item["mention"] for item in group] most_common_raw = Counter(raw_names).most_common(1)[0][0] acronyms = [clean_acronym(item["acronym"]).lower() for item in group if item["acronym"]] best_acronym = Counter(acronyms).most_common(1)[0][0] if acronyms else "" canonical_name = None for lookup_key in [cleaned, clean_name(most_common_raw), best_acronym.lower()]: if lookup_key and lookup_key in harmonization_map: canonical_name = harmonization_map[lookup_key] harmonization_hits += 1 break if canonical_name is None: if cleaned in dynamic_acronym_map: canonical_name = dynamic_acronym_map[cleaned] fallback_hits += 1 elif best_acronym and best_acronym in dynamic_acronym_map: canonical_name = dynamic_acronym_map[best_acronym] fallback_hits += 1 if canonical_name is None: canonical_name = most_common_raw for item in group: raw_to_canonical[item["mention"]] = canonical_name print(f" Harmonization map resolved: {harmonization_hits} groups") print(f" Fallback acronym map resolved: {fallback_hits} groups") # Apply canonical mapping to named mentions only for m in named_mentions: m["canonical"] = raw_to_canonical[m["mention"]] def standardize_typology(typ: str) -> str: typ = typ.lower() if "survey" in typ or "microdata" in typ: return "Survey" elif "admin" in typ or "registry" in typ or "records" in typ: return "Administrative Data" elif "census" in typ: return "Census" elif "satellite" in typ or "spatial" in typ or "remote" in typ: return "Geospatial/Satellite" elif "indicators" in typ or "macro" in typ or "aggregate" in typ: return "Macro/Aggregate Indicators" elif not typ: return "Unknown" else: return typ.title() for m in named_mentions: m["typology"] = standardize_typology(m["typology"]) for m in descriptive_mentions: m["typology"] = standardize_typology(m["typology"]) for m in vague_mentions: m["typology"] = standardize_typology(m["typology"]) print("Step 3: Calculating frequencies and trends...") # Calculate frequencies per canonical dataset dataset_papers = defaultdict(set) dataset_mentions_count = Counter() dataset_attributes = defaultdict(lambda: { "acronyms": Counter(), "producers": Counter(), "geographies": Counter(), "typologies": Counter(), }) # Group papers details by dataset paper_lookup = {p["id"]: p for p in papers} for m in named_mentions: canonical = m["canonical"] paper_id = m["paper_id"] dataset_papers[canonical].add(paper_id) dataset_mentions_count[canonical] += 1 if m["acronym"]: dataset_attributes[canonical]["acronyms"][m["acronym"]] += 1 if m["producer"]: dataset_attributes[canonical]["producers"][m["producer"]] += 1 if m["geography"]: dataset_attributes[canonical]["geographies"][m["geography"]] += 1 if m["typology"]: dataset_attributes[canonical]["typologies"][m["typology"]] += 1 # Pre-group named mentions by canonical for O(M) paper-page lookup mentions_by_canonical = defaultdict(list) for m in named_mentions: mentions_by_canonical[m["canonical"]].append(m) # Format the top datasets list top_datasets_list = [] for canonical, paper_ids in dataset_papers.items(): df = len(paper_ids) mf = dataset_mentions_count[canonical] attrs = dataset_attributes[canonical] best_acronym = attrs["acronyms"].most_common(1)[0][0] if attrs["acronyms"] else "" best_producer = attrs["producers"].most_common(1)[0][0] if attrs["producers"] else "Unknown" best_geography = attrs["geographies"].most_common(1)[0][0] if attrs["geographies"] else "Global / Multiple" best_typology = attrs["typologies"].most_common(1)[0][0] if attrs["typologies"] else "Unknown" # Collect paper references (uses pre-grouped O(M) lists) referencing_papers = [] paper_pages = defaultdict(list) for m in mentions_by_canonical[canonical]: paper_pages[m["paper_id"]].append(m["page"]) for p_id in paper_ids: p_info = paper_lookup[p_id] referencing_papers.append({ "id": p_id, "title": p_info["title"], "year": p_info["year"], "authors": p_info["authors"], "pdf_url": p_info.get("pdf_url", ""), "pages": sorted(list(set(paper_pages[p_id]))) }) referencing_papers.sort(key=lambda x: x["year"] if isinstance(x["year"], int) else 0, reverse=True) # Collect all raw variants and their frequencies that resolved to this canonical variant_counter = Counter(m["mention"] for m in mentions_by_canonical[canonical]) sorted_variants = [{"name": name, "count": count} for name, count in variant_counter.most_common()] top_datasets_list.append({ "canonical_name": canonical, "acronym": best_acronym, "typology": best_typology, "producer": best_producer, "geography": best_geography, "document_frequency": df, "mention_frequency": mf, "variants": sorted_variants, "papers": referencing_papers }) top_datasets_list.sort(key=lambda x: x["document_frequency"], reverse=True) # Calculate Annual Trends papers_by_year = defaultdict(list) for p in papers: if isinstance(p["year"], int): papers_by_year[p["year"]].append(p) annual_trends = [] valid_years = sorted([y for y in papers_by_year.keys() if 2000 <= y <= 2026]) for y in valid_years: year_papers = papers_by_year[y] total = len(year_papers) with_data = sum(1 for p in year_papers if p["has_data"]) pct = (with_data / total * 100) if total > 0 else 0 annual_trends.append({ "year": y, "total_papers": total, "papers_with_data": with_data, "percentage": round(pct, 1) }) # ── Named dataset distributions (for top-datasets chart and typology donut) ── typology_counts = Counter() producer_counts = Counter() geography_counts = Counter() for ds in top_datasets_list: df = ds["document_frequency"] typology_counts[ds["typology"]] += df if ds["producer"] != "Unknown": producer_counts[ds["producer"]] += df if ds["geography"] != "Global / Multiple": geography_counts[ds["geography"]] += df # ── Data Practice Signal breakdown (descriptive and vague) ── practice_papers = defaultdict(set) practice_mentions = Counter() for m in descriptive_mentions + vague_mentions: spec = m["specificity"].lower() key = (m["typology"], spec) practice_papers[key].add(m["paper_id"]) practice_mentions[key] += 1 descriptive_breakdown = [ { "typology": key[0], "specificity": key[1].upper(), "document_frequency": len(paper_ids), "mention_frequency": practice_mentions[key], } for key, paper_ids in sorted( practice_papers.items(), key=lambda x: len(x[1]), reverse=True ) ] papers_with_named = len({m["paper_id"] for m in named_mentions}) papers_with_descriptive_only = len( {m["paper_id"] for m in descriptive_mentions} - {m["paper_id"] for m in named_mentions} ) # ── Specificity, Usage Context, and Is Used distributions ── all_mentions = named_mentions + descriptive_mentions + vague_mentions specificity_counts = Counter() usage_context_counts = Counter() is_used_counts = Counter() for m in all_mentions: # 1. Specificity (Named vs Descriptive vs Vague) spec = m.get("specificity", "unknown").strip().lower() if spec == "named": specificity_counts["Named"] += 1 elif spec == "descriptive": specificity_counts["Descriptive"] += 1 elif spec == "vague": specificity_counts["Vague"] += 1 else: specificity_counts[spec.capitalize()] += 1 # 2. Usage Context ctx = m.get("usage_context") if not ctx: usage_context_counts["Unknown"] += 1 elif ctx.lower() == "primary": usage_context_counts["Primary Use"] += 1 elif ctx.lower() == "supporting": usage_context_counts["Supporting Use"] += 1 elif ctx.lower() == "background": usage_context_counts["Background / Citation"] += 1 else: usage_context_counts[ctx.replace("_", " ").title()] += 1 # 3. Is Used used_val = m.get("is_used") if not used_val: is_used_counts["Unknown"] += 1 elif used_val == "True": is_used_counts["Used in Analysis"] += 1 elif used_val == "False": is_used_counts["Not Used (Citation Only)"] += 1 else: is_used_counts[used_val.replace("_", " ").title()] += 1 # ── Build final dashboard payload ── papers_with_data = sum(1 for p in papers if p["has_data"]) payload = { "summary": { "total_papers": len(papers), "papers_with_data": papers_with_data, "data_adoption_rate": round(papers_with_data / len(papers) * 100, 1), # Named-only counts (for dataset entity explorer) "named_mentions": len(named_mentions), "unique_canonical_datasets": len(top_datasets_list), "papers_with_named_datasets": papers_with_named, # Descriptive-only counts (for data practice section) "descriptive_mentions": len(descriptive_mentions), "papers_with_descriptive_only": papers_with_descriptive_only, # Vague counts "vague_mentions": len(vague_mentions), }, "annual_trends": annual_trends, # Named dataset typology distribution "typology_distribution": [ {"name": name, "value": val} for name, val in typology_counts.most_common(10) ], "specificity_distribution": [ {"name": name, "value": val} for name, val in specificity_counts.most_common() ], "usage_context_distribution": [ {"name": name, "value": val} for name, val in usage_context_counts.most_common() ], "is_used_distribution": [ {"name": name, "value": val} for name, val in is_used_counts.most_common() ], "top_producers": [ {"name": name, "value": val} for name, val in producer_counts.most_common(10) ], "top_geographies": [ {"name": name, "value": val} for name, val in geography_counts.most_common(10) ], # Named dataset entities (entity explorer + knowledge graph) "datasets": top_datasets_list, # Descriptive mention breakdown (data-practice signal) "descriptive_breakdown": descriptive_breakdown, } # Save to data/dashboard_data.json print(f"Step 4: Writing output dashboard JSON and JS...") os.makedirs(os.path.dirname(output_json_path), exist_ok=True) with open(output_json_path, "w", encoding="utf-8") as f: json.dump(payload, f, indent=2, ensure_ascii=False) print(f" Saved dashboard JSON to: {output_json_path}") # Also save to data/dashboard_data.js for visual dashboard import output_js_path = base_dir / "data" / "dashboard_data.js" with open(output_js_path, "w", encoding="utf-8") as f: f.write("const DASHBOARD_DATA = ") json.dump(payload, f, ensure_ascii=False) f.write(";") print(f" Saved dashboard JS to: {output_js_path}") # Export flat deduplicated mentions CSV print(f"Step 5: Exporting deduplicated datasets CSV...") with open(output_csv_path, "w", newline="", encoding="utf-8") as f: writer = csv.writer(f) writer.writerow([ "canonical_name", "acronym", "typology", "producer", "geography", "document_frequency", "mention_frequency" ]) for ds in top_datasets_list: writer.writerow([ ds["canonical_name"], ds["acronym"], ds["typology"], ds["producer"], ds["geography"], ds["document_frequency"], ds["mention_frequency"] ]) print(f" Saved CSV to: {output_csv_path}") print("Step 6: Generating Graph Database Neo4j Import CSVs...") os.makedirs(graph_base, exist_ok=True) # 6.1 Nodes: Papers with open(graph_base / "nodes_papers.csv", "w", newline="", encoding="utf-8") as f: writer = csv.writer(f) writer.writerow(["id", "title", "year"]) for p in papers: writer.writerow([p["id"], p["title"], p["year"]]) # 6.2 Nodes: Datasets # Assign each canonical dataset a unique slug/ID dataset_slugs = {} with open(graph_base / "nodes_datasets.csv", "w", newline="", encoding="utf-8") as f: writer = csv.writer(f) writer.writerow(["id", "name", "acronym", "typology", "producer", "geography"]) for i, ds in enumerate(top_datasets_list): slug = sanitize_author_id(ds["canonical_name"]) # Ensure unique ID if slug in dataset_slugs.values(): slug = f"{slug}_{i}" dataset_slugs[ds["canonical_name"]] = slug writer.writerow([ slug, ds["canonical_name"], ds["acronym"], ds["typology"], ds["producer"], ds["geography"] ]) # 6.3 Nodes: Authors author_ids = {} with open(graph_base / "nodes_authors.csv", "w", newline="", encoding="utf-8") as f: writer = csv.writer(f) writer.writerow(["id", "name"]) for i, author in enumerate(sorted(all_authors)): a_id = sanitize_author_id(author) if a_id in author_ids.values(): a_id = f"{a_id}_{i}" author_ids[author] = a_id writer.writerow([a_id, author]) # 6.4 Edges: Authored (Author -> Paper) with open(graph_base / "edges_authored.csv", "w", newline="", encoding="utf-8") as f: writer = csv.writer(f) writer.writerow(["author_id", "paper_id"]) for paper_id, authors_list in paper_to_authors.items(): for author in authors_list: writer.writerow([author_ids[author], paper_id]) # 6.5 Edges: Mentions (Paper -> Dataset) # Collect relationship properties # A paper can mention a dataset multiple times on different pages paper_dataset_edges = defaultdict(lambda: { "pages": set(), "contexts": set(), "confidences": [] }) for m in named_mentions: edge_key = (m["paper_id"], m["canonical"]) paper_dataset_edges[edge_key]["pages"].add(m["page"]) if m["usage_context"]: paper_dataset_edges[edge_key]["contexts"].add(m["usage_context"]) paper_dataset_edges[edge_key]["confidences"].append(m["confidence"]) with open(graph_base / "edges_mentions.csv", "w", newline="", encoding="utf-8") as f: writer = csv.writer(f) writer.writerow(["paper_id", "dataset_id", "pages", "context", "confidence"]) for (paper_id, canonical), data in paper_dataset_edges.items(): dataset_id = dataset_slugs[canonical] pages_str = ";".join(map(str, sorted(list(data["pages"])))) context = "|".join(data["contexts"]) if data["contexts"] else "Unknown" avg_conf = round(sum(data["confidences"]) / len(data["confidences"]), 3) if data["confidences"] else 0.0 writer.writerow([paper_id, dataset_id, pages_str, context, avg_conf]) # Create README_GRAPH.md print("Step 7: Creating Cypher Import Guide...") with open(graph_base / "README_GRAPH.md", "w", encoding="utf-8") as f: f.write("""# Neo4j Graph Database Import Guide This folder contains Neo4j-import-ready CSV files containing: * **Nodes**: Papers, Authors, and deduplicated Datasets. * **Relationships**: Authorship (`:AUTHORED`) and dataset citations (`:MENTIONS`). --- ## CSV File List 1. **`nodes_papers.csv`**: Contains Policy Research Working Papers. 2. **`nodes_datasets.csv`**: Contains canonicalized datasets. 3. **`nodes_authors.csv`**: Contains unique authors. 4. **`edges_authored.csv`**: Maps Authors to Papers. 5. **`edges_mentions.csv`**: Maps Papers to Datasets with page numbers, context, and model confidence scores. --- ## Import Cypher Queries To import these files into your Neo4j instance, place the CSV files in your Neo4j project's `import/` directory, open the Neo4j Browser, and execute the following queries: ### 1. Create Constraints ```cypher CREATE CONSTRAINT UNIQUE_paper FOR (p:Paper) REQUIRE p.id IS UNIQUE; CREATE CONSTRAINT UNIQUE_dataset FOR (d:Dataset) REQUIRE d.id IS UNIQUE; CREATE CONSTRAINT UNIQUE_author FOR (a:Author) REQUIRE a.id IS UNIQUE; ``` ### 2. Load Nodes ```cypher // Load Papers LOAD CSV WITH HEADERS FROM 'file:///nodes_papers.csv' AS row MERGE (p:Paper {id: row.id}) SET p.title = row.title, p.year = toInteger(row.year); // Load Datasets LOAD CSV WITH HEADERS FROM 'file:///nodes_datasets.csv' AS row MERGE (d:Dataset {id: row.id}) SET d.name = row.name, d.acronym = row.acronym, d.typology = row.typology, d.producer = row.producer, d.geography = row.geography; // Load Authors LOAD CSV WITH HEADERS FROM 'file:///nodes_authors.csv' AS row MERGE (a:Author {id: row.id}) SET a.name = row.name; ``` ### 3. Load Relationships ```cypher // Load AUTHORED LOAD CSV WITH HEADERS FROM 'file:///edges_authored.csv' AS row MATCH (a:Author {id: row.author_id}) MATCH (p:Paper {id: row.paper_id}) MERGE (a)-[:AUTHORED]->(p); // Load MENTIONS LOAD CSV WITH HEADERS FROM 'file:///edges_mentions.csv' AS row MATCH (p:Paper {id: row.paper_id}) MATCH (d:Dataset {id: row.dataset_id}) MERGE (p)-[:MENTIONS { pages: split(row.pages, ';'), context: row.context, confidence: toFloat(row.confidence) }]->(d); ``` """) print(f" Saved Cypher guide to: {graph_base / 'README_GRAPH.md'}") print("\nAll pipeline files generated successfully!") if __name__ == "__main__": main()