From 2f4326ec151dc88b15aa735973960b0c99dfb96f Mon Sep 17 00:00:00 2001 From: NguyenND Date: Fri, 13 Feb 2026 17:24:16 -0500 Subject: [PATCH] Update import function --- app.py | 197 +++++++++++++++++++++++++++++++++++++++++++++++++-------- 1 file changed, 171 insertions(+), 26 deletions(-) diff --git a/app.py b/app.py index 395a036..7d36a57 100644 --- a/app.py +++ b/app.py @@ -112,7 +112,13 @@ def inject_common_variables(): def log_action(action: str, details: str = "", level: str = "info"): """Log user actions for audit trail.""" timestamp = datetime.now().isoformat() - user_ip = request.remote_addr + + # Handle case when called outside request context (e.g., background threads) + try: + user_ip = request.remote_addr + except RuntimeError: + user_ip = "background" + log_message = f"[{timestamp}] [{user_ip}] [{action}] {details}" if level == "error": @@ -727,14 +733,76 @@ def import_data(): }), 400 -# Import job storage (in-memory for simplicity) -import_jobs = {} +# ============================================================================= +# Import Job Storage (File-based for multi-worker support) +# ============================================================================= +import json +import threading + +JOBS_DIR = os.path.join(os.path.dirname(__file__), 'data', 'jobs') +os.makedirs(JOBS_DIR, exist_ok=True) + +# Thread lock for file operations +job_lock = threading.Lock() + +def save_job(job_id, job_data): + """Save job data to file.""" + with job_lock: + job_file = os.path.join(JOBS_DIR, f'{job_id}.json') + # Convert datetime to string for JSON serialization + job_copy = job_data.copy() + if 'start_time' in job_copy and hasattr(job_copy['start_time'], 'isoformat'): + job_copy['start_time'] = job_copy['start_time'].isoformat() + with open(job_file, 'w') as f: + json.dump(job_copy, f) + +def load_job(job_id): + """Load job data from file.""" + job_file = os.path.join(JOBS_DIR, f'{job_id}.json') + if not os.path.exists(job_file): + return None + try: + with open(job_file, 'r') as f: + job_data = json.load(f) + # Convert start_time back to datetime if needed + if 'start_time' in job_data and isinstance(job_data['start_time'], str): + job_data['start_time'] = datetime.fromisoformat(job_data['start_time']) + return job_data + except Exception as e: + logger.error(f"Error loading job {job_id}: {e}") + return None + +def update_job(job_id, updates): + """Update specific fields in a job.""" + with job_lock: + job_data = load_job(job_id) + if job_data: + job_data.update(updates) + save_job(job_id, job_data) + return job_data + return None + +def cleanup_old_jobs(): + """Remove job files older than 1 hour.""" + try: + now = datetime.now() + for filename in os.listdir(JOBS_DIR): + if filename.endswith('.json'): + filepath = os.path.join(JOBS_DIR, filename) + file_time = datetime.fromtimestamp(os.path.getmtime(filepath)) + if (now - file_time).total_seconds() > 3600: # 1 hour + os.remove(filepath) + except Exception as e: + logger.warning(f"Error cleaning up old jobs: {e}") + @app.route('/api/import/start', methods=['POST']) def start_import(): """Start an async import job.""" import uuid - import threading + + # Cleanup old jobs + cleanup_old_jobs() try: data = request.json @@ -769,7 +837,7 @@ def start_import(): # Create job job_id = str(uuid.uuid4())[:8] - import_jobs[job_id] = { + job_data = { 'status': 'running', 'total_records': parse_result.total_rows, 'processed': 0, @@ -782,58 +850,118 @@ def start_import(): 'error': None, 'duration_seconds': 0 } + save_job(job_id, job_data) # Store credentials for thread creds_data = session.get('qbo_credentials') + if not creds_data: + # Get from settings if not in session + creds = settings.get_credentials() + creds_data = asdict(creds) if creds else None + import_settings_data = asdict(settings.import_settings) + # Serialize parse_result for thread + parse_result_data = { + 'filepath': parse_result.filepath, + 'sheet_name': parse_result.sheet_name, + 'columns': parse_result.columns, + 'total_rows': parse_result.total_rows, + 'valid_rows': parse_result.valid_rows, + 'error_count': parse_result.error_count, + 'warning_count': parse_result.warning_count, + 'rows': [] + } + for row in parse_result.rows: + parse_result_data['rows'].append({ + 'row_number': row.row_number, + 'original_data': row.original_data, + 'qbo_data': row.qbo_data, + 'is_valid': row.is_valid + }) + # Start background import def run_import(): try: from src.config.settings import QBOCredentials, ImportSettings + # Load job from file + job = load_job(job_id) + if not job: + logger.error(f"Job {job_id} not found in run_import") + return + # Recreate client in thread creds = QBOCredentials(**creds_data) if creds_data else settings.get_credentials() thread_client = QBOClient(creds) thread_settings = ImportSettings(**import_settings_data) - job = import_jobs[job_id] - # Create processor processor = ImportProcessor(thread_client, thread_settings) + # Rebuild parse result for validation + from src.core.excel_parser import ParsedRow, ParseResult + rebuilt_rows = [] + for rd in parse_result_data['rows']: + row = ParsedRow( + row_number=rd['row_number'], + original_data=rd['original_data'], + qbo_data=rd['qbo_data'], + is_valid=rd['is_valid'] + ) + rebuilt_rows.append(row) + + rebuilt_parse_result = ParseResult( + filepath=parse_result_data['filepath'], + sheet_name=parse_result_data['sheet_name'], + columns=parse_result_data['columns'], + rows=rebuilt_rows, + total_rows=parse_result_data['total_rows'], + valid_rows=parse_result_data['valid_rows'], + error_count=parse_result_data['error_count'], + warning_count=parse_result_data['warning_count'] + ) + # Validate and resolve records - records = processor.validate_and_resolve(parse_result, data_type) + records = processor.validate_and_resolve(rebuilt_parse_result, data_type) # Process each record from src.core.import_processor import ImportStatus for i, record in enumerate(records): + # Reload job to get latest state + job = load_job(job_id) + if not job: + logger.error(f"Job {job_id} disappeared during import") + return + + result_entry = None + if record.status == ImportStatus.FAILED: job['failed'] += 1 - job['recent_results'].append({ + result_entry = { 'row': record.row_number, 'status': 'failed', 'error': record.error_message, 'key_fields': None, 'data_type': data_type - }) + } elif record.status == ImportStatus.DUPLICATE: job['duplicates'] += 1 - job['recent_results'].append({ + result_entry = { 'row': record.row_number, 'status': 'duplicate', 'error': None, 'data_type': data_type - }) + } elif record.status == ImportStatus.SKIPPED: job['skipped'] += 1 - job['recent_results'].append({ + result_entry = { 'row': record.row_number, 'status': 'skipped', 'error': record.error_message, 'data_type': data_type - }) + } elif record.status == ImportStatus.PENDING: # Import to QBO try: @@ -856,47 +984,64 @@ def start_import(): qbo_id = qbo_entity.get("Id") job['successful'] += 1 - job['recent_results'].append({ + result_entry = { 'row': record.row_number, 'status': 'success', 'qbo_id': qbo_id, 'data_type': data_type - }) + } except Exception as e: job['failed'] += 1 key_fields = processor._get_key_fields(record.qbo_data, data_type) - job['recent_results'].append({ + result_entry = { 'row': record.row_number, 'status': 'failed', 'error': str(e), 'key_fields': key_fields, 'data_type': data_type - }) + } # Small delay to avoid rate limiting import time time.sleep(0.1) + # Update job progress + if result_entry: + job['recent_results'].append(result_entry) job['processed'] = i + 1 - # Keep only last 50 results to limit memory + # Keep only last 50 results to limit file size if len(job['recent_results']) > 50: job['recent_results'] = job['recent_results'][-50:] + + # Save job state + save_job(job_id, job) - job['status'] = 'completed' - job['duration_seconds'] = (datetime.now() - job['start_time']).total_seconds() + # Final update + job = load_job(job_id) + if job: + job['status'] = 'completed' + start_time = job.get('start_time') + if isinstance(start_time, str): + start_time = datetime.fromisoformat(start_time) + job['duration_seconds'] = (datetime.now() - start_time).total_seconds() + save_job(job_id, job) log_action("CREATE", f"Import complete: {job['successful']} successful, {job['failed']} failed") except Exception as e: import traceback - job = import_jobs.get(job_id) + logger.error(f"Import thread error: {traceback.format_exc()}") + job = load_job(job_id) if job: job['status'] = 'failed' job['error'] = str(e) - job['duration_seconds'] = (datetime.now() - job['start_time']).total_seconds() - logger.error(f"Import thread error: {traceback.format_exc()}") + start_time = job.get('start_time') + if isinstance(start_time, str): + start_time = datetime.fromisoformat(start_time) + job['duration_seconds'] = (datetime.now() - start_time).total_seconds() + save_job(job_id, job) thread = threading.Thread(target=run_import) thread.daemon = True @@ -922,7 +1067,7 @@ def start_import(): @app.route('/api/import/progress/', methods=['GET']) def get_import_progress(job_id): """Get import job progress.""" - job = import_jobs.get(job_id) + job = load_job(job_id) if not job: return jsonify({'success': False, 'error': 'Job not found'}), 404 @@ -943,7 +1088,7 @@ def get_import_progress(job_id): @app.route('/api/import/result/', methods=['GET']) def get_import_result(job_id): """Get final import result.""" - job = import_jobs.get(job_id) + job = load_job(job_id) if not job: return jsonify({'success': False, 'error': 'Job not found'}), 404