Update import function
This commit is contained in:
@@ -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:]
|
||||
|
||||
job['status'] = 'completed'
|
||||
job['duration_seconds'] = (datetime.now() - job['start_time']).total_seconds()
|
||||
# Save job state
|
||||
save_job(job_id, job)
|
||||
|
||||
# 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/<job_id>', 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/<job_id>', 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
|
||||
|
||||
|
||||
Reference in New Issue
Block a user