#!/usr/bin/env python3 """ PLATFORM-4936 pipeline — Step 1: transform source CSV → bulk_runner input. Reads the business-supplied PLATFORM-4946_cleaned_up_users.csv, looks up staff uuids in neo4j and removes those which have already been reactivated to avoid interfering with any active users. Only users where the auth0UserId is the same as the id should be revoked. Usage: uv run python src/pipelines/platform-4946/1_prepare.py \\ --input "/data/prepare/PLATFORM-4946_cleaned_up_users.csv" \\ --output data/output/platform-4946_identities.csv Neo4j creds are read from .env (NEO4J_*). Requires awsume prod. """ import argparse import csv import os import sys from pathlib import Path from typing import Generator import connector_neo4j from dotenv import load_dotenv load_dotenv() NEO4J_URL = os.environ.get('NEO4J_URL') NEO4J_USERNAME = os.environ.get('NEO4J_USERNAME') NEO4J_PASSWORD = os.environ.get('NEO4J_PASSWORD') NEO4J_MAX_RETRY_TIME = os.environ.get('NEO4J_MAX_RETRY_TIME') or 60 connector_neo4j.configure( NEO4J_URL, NEO4J_USERNAME, NEO4J_PASSWORD, max_transaction_retry_time=NEO4J_MAX_RETRY_TIME, ) CYPHER_QUERY_IS_INACTIVE = """ WITH $id_list AS identities MATCH (i:Identity) WHERE i.id in identities AND i.auth0UserId = i.id AND i.active = "N" RETURN i.id as id, i.email as email """ def get_inactive(identity_ids: list[str]) -> Generator[tuple[str, str]]: session: connector_neo4j.Session = connector_neo4j.get_session() result = session.run(CYPHER_QUERY_IS_INACTIVE, id_list=identity_ids) return ((identity['id'], identity['email']) for identity in result) @connector_neo4j.Neo4jSession(transaction=True, use_v2=True, database='graph.db') def main() -> None: parser = argparse.ArgumentParser(description='PLATFORM-4946 Step 1: remove active users and format CSV') parser.add_argument('--input', required=True, help='Source CSV path') parser.add_argument('--output', required=True, help='Output CSV path') args = parser.parse_args() input_path = Path(args.input) output_path = Path(args.output) if not input_path.exists(): print(f'Input not found: {input_path}', file=sys.stderr) sys.exit(1) # Minus 1 removes the heading line raw_count = sum(1 for _ in input_path.open('rb')) - 1 with input_path.open() as infile: identity_ids = [row['identity_id'] for row in csv.DictReader(infile) if row['status'] == 'success'] filtered_count = 0 with output_path.open(mode='w') as outfile: outwriter = csv.writer(outfile) outwriter.writerow(('id', 'email')) for row in get_inactive(identity_ids): outwriter.writerow(row) filtered_count += 1 print(f'Wrote {output_path} with {filtered_count} ids filtered from original {raw_count} identities.') if __name__ == '__main__': main()