#!/usr/bin/env python3 """ Create vendor accounts from CSV using the ows-account API. Features: - Automatic retry for transient failures - Progress saving and resume capability - Graceful interrupt handling (Ctrl+C) - Account ID output CSV generation """ import argparse import csv import json import logging import os import sys import time from typing import Dict, List, Optional from helpers import ( load_config, parse_csv_row, create_vendor, write_success_csv, write_failures_csv, print_summary ) # Configure logging logging.basicConfig( level=logging.INFO, format='%(asctime)s - %(levelname)s - %(message)s' ) logger = logging.getLogger(__name__) def process_csv( csv_file: str, api_url: str, bearer_tokens: Dict[str, str], service_tier_map: Dict[str, str], identity_headers: Optional[Dict[str, str]] = None, dry_run: bool = False, limit: Optional[int] = None, delay: float = 0.5, default_owner: str = None, default_service_tier: str = None, skip_if_missing: bool = False, output_file: str = None, success_csv: str = None, skip_existing: bool = False, append: bool = False ) -> Dict[str, List]: """Process CSV file and create vendors with progress saving.""" results = { 'success': [], 'errors': [], 'skipped': [] } success_rows = [] # Only contains NEW rows to be written processed_identifiers = set() if success_csv and os.path.exists(success_csv) and os.path.getsize(success_csv) > 0: try: with open(success_csv, 'r', encoding='utf-8') as f: existing_reader = csv.DictReader(f) for existing_row in existing_reader: # Only track identifiers for skipping, don't add to success_rows # Use composite key: Royalty Share ID + Account Name royalty_id = existing_row.get('RoyaltyShare ID', '').strip() account_name = existing_row.get('Account Name', '').strip() identifier = f"{royalty_id}|{account_name}" if (royalty_id or account_name) else None if identifier: processed_identifiers.add(identifier) logger.info(f"[RESUME] Found existing results with {len(processed_identifiers)} accounts - will skip and append") skip_existing = True append = True except Exception as e: logger.warning(f"Could not load existing success CSV: {e}") elif skip_existing: logger.info(f"[SKIP] Skip-existing mode enabled but no existing results found") save_interval = 10 original_fieldnames = None input_file_seen = set() # Track identifiers seen in current input file try: with open(csv_file, 'r', encoding='utf-8') as f: reader = csv.DictReader(f) original_fieldnames = reader.fieldnames for idx, row in enumerate(reader, start=1): if limit and idx > limit: logger.info(f"Reached limit of {limit} rows") break # Use composite key: Royalty Share ID + Account Name royalty_id = row.get('RoyaltyShare ID', '').strip() account_name = row.get('Account Name', '').strip() identifier = f"{royalty_id}|{account_name}" if (royalty_id or account_name) else None # Skip if duplicate within input file if identifier in input_file_seen: logger.warning(f"[DUPLICATE] Row {idx} is duplicate in input file: {account_name} (Royalty ID: {royalty_id})") results['skipped'].append({ 'row': idx, 'data': row, 'reason': f'Duplicate in input file (Royalty ID: {royalty_id}, Name: {account_name})' }) continue # Skip if already processed in previous runs if skip_existing and identifier in processed_identifiers: logger.info(f"[SKIP] Row {idx} already processed: {row.get('Account Name', 'Unknown')}") continue # Mark as seen in current input if identifier: input_file_seen.add(identifier) logger.info(f"Processing row {idx}: {row.get('Account Name', 'Unknown')}") payload = parse_csv_row( row, service_tier_map=service_tier_map, default_owner=default_owner, default_service_tier=default_service_tier, skip_if_missing=skip_if_missing ) if not payload: results['skipped'].append({ 'row': idx, 'data': row, 'reason': 'Invalid data' }) continue company_brand = payload['company_brand'] bearer_token = bearer_tokens.get(company_brand) or bearer_tokens.get('default') if not bearer_token: logger.error(f"No bearer token configured for company brand: {company_brand}") results['skipped'].append({ 'row': idx, 'payload': payload, 'reason': f'No bearer token for {company_brand}' }) continue result = create_vendor(api_url, payload, bearer_token, identity_headers, dry_run) if result['status'] == 'success' or result['status'] == 'dry_run': results['success'].append(result) success_row_data = dict(row) if result['status'] == 'success' and 'data' in result: account_id = result['data'].get('vendor_id', result['data'].get('id', 'N/A')) success_row_data['Account ID'] = account_id elif result['status'] == 'dry_run': success_row_data['Account ID'] = 'DRY_RUN' success_rows.append(success_row_data) # Track as processed for auto-resume if identifier: processed_identifiers.add(identifier) else: results['errors'].append({'row': idx, **result}) if not dry_run: time.sleep(delay) if output_file and idx % save_interval == 0: try: with open(output_file, 'w') as f: json.dump(results, f, indent=2) logger.info(f"[PROGRESS] Progress saved ({idx} rows processed)") except Exception as e: logger.warning(f"Failed to save progress: {e}") if success_csv and success_rows and idx % save_interval == 0: try: write_success_csv(success_rows, success_csv, original_fieldnames, append=append) logger.info(f"[PROGRESS] Success CSV updated ({len(success_rows)} new)") success_rows = [] # Clear after saving to prevent duplicates append = True # Ensure subsequent saves append except Exception as e: logger.warning(f"Failed to save success CSV: {e}") except FileNotFoundError: logger.error(f"CSV file not found: {csv_file}") sys.exit(1) except KeyboardInterrupt: logger.warning("\n[INTERRUPTED] Interrupted by user. Saving progress...") if output_file: try: with open(output_file, 'w') as f: json.dump(results, f, indent=2) logger.info(f"[SAVED] Progress saved to: {output_file}") except Exception as e: logger.error(f"Failed to save progress: {e}") if success_csv and success_rows: try: write_success_csv(success_rows, success_csv, original_fieldnames, append=append) logger.info(f"[SAVED] Success CSV saved to: {success_csv}") except Exception as e: logger.error(f"Failed to save success CSV: {e}") sys.exit(130) except Exception as e: logger.error(f"Unexpected error processing CSV: {e}") if output_file: try: with open(output_file, 'w') as f: json.dump(results, f, indent=2) logger.warning(f"[SAVED] Partial results saved to: {output_file}") except Exception: pass if success_csv and success_rows: try: write_success_csv(success_rows, success_csv, original_fieldnames, append=append) logger.warning(f"[SAVED] Partial success CSV saved to: {success_csv}") except Exception: pass sys.exit(1) # Write final success CSV if success_csv and success_rows: write_success_csv(success_rows, success_csv, original_fieldnames, append=append) action = "Appended" if append else "Wrote" logger.info(f"[CSV] {action} {len(success_rows)} new accounts to: {success_csv}") # Write failures CSV if (results['errors'] or results['skipped']) and success_csv: failures_csv = success_csv.replace('.csv', '_failures.csv') write_failures_csv(results, failures_csv, original_fieldnames) logger.info(f"[CSV] Wrote {len(results['errors']) + len(results['skipped'])} failures to: {failures_csv}") return results def main(): parser = argparse.ArgumentParser( description='Create vendor accounts from CSV using v2_create_vendor endpoint' ) parser.add_argument( '--csv', required=True, help='Path to CSV file with account data' ) parser.add_argument( '--config', required=True, help='Path to JSON config file with bearer tokens (see config.example.json)' ) parser.add_argument( '--env', choices=['qa', 'prod'], default='qa', help='Environment: qa or prod (default: qa)' ) parser.add_argument( '--api-url', help='Base API URL (overrides --env if specified)' ) parser.add_argument( '--dry-run', action='store_true', help='Dry run - show what would be created without making API calls' ) parser.add_argument( '--limit', type=int, help='Limit number of rows to process (for testing)' ) parser.add_argument( '--delay', type=float, default=0.5, help='Delay in seconds between API calls (default: 0.5)' ) parser.add_argument( '--output', help='Output JSON file for results' ) parser.add_argument( '--default-owner', help='Default owner to use when CSV Owner field is empty (e.g., "odd")' ) parser.add_argument( '--default-service-tier', help='Default service tier UUID to use when CSV Service Tier is empty (e.g., 5f2bd4fc-df94-4f35-97d3-ef23f8573279 for diy-tier-1)' ) parser.add_argument( '--skip-if-missing', action='store_true', help='Skip rows with missing required fields instead of using defaults' ) parser.add_argument( '--success-csv', help='Output CSV file for successful account creations with Account ID prepended' ) parser.add_argument( '--skip-existing', action='store_true', help='Skip rows that already exist in the success CSV (auto-enabled when success CSV exists)' ) parser.add_argument( '--append', action='store_true', help='Append results to existing success CSV (auto-enabled when success CSV exists)' ) args = parser.parse_args() # Determine API URL if args.api_url: api_url = args.api_url else: api_url = 'https://ows-account.theorchard.io' if args.env == 'prod' else 'https://ows-grass.theorchard.io' # Load configuration logger.info(f"Loading config from: {args.config}") config = load_config(args.config) # Process CSV logger.info(f"Processing CSV: {args.csv}") logger.info(f"Environment: {args.env.upper()}") logger.info(f"API URL: {api_url}") if args.dry_run: logger.info("DRY RUN MODE - No API calls will be made") results = process_csv( csv_file=args.csv, api_url=api_url, bearer_tokens=config['bearer_token'], service_tier_map=config['service_tier_map'], identity_headers=config.get('identity_headers'), # Optional dry_run=args.dry_run, limit=args.limit, delay=args.delay, default_owner=args.default_owner, default_service_tier=args.default_service_tier, skip_if_missing=args.skip_if_missing, output_file=args.output, success_csv=args.success_csv, skip_existing=args.skip_existing, append=args.append ) # Print summary and save results print_summary(results) if args.output: with open(args.output, 'w') as f: json.dump(results, f, indent=2) logger.info(f"\nResults saved to: {args.output}") if __name__ == '__main__': main()