#!/usr/bin/env python3 """ Reassemble split Datadog logs that were fragmented by fluentbit/awsfirelens. Groups log fragments by partial_id and orders them by partial_ordinal. """ import csv import json import os import sys from pathlib import Path from collections import defaultdict def reassemble_split_logs(csv_file_path, output_dir): """ Read CSV file with split logs and reassemble them. Returns: tuple: (complete_logs, incomplete_logs, stats) """ # Track log fragments grouped by partial_id log_fragments = defaultdict(list) complete_logs = [] with open(csv_file_path, 'r', encoding='utf-8') as csvfile: reader = csv.DictReader(csvfile) # Validate required columns exist required_columns = ['partial_id', 'partial_ordinal', 'partial_last'] if reader.fieldnames is None: raise ValueError("CSV file appears to be empty or invalid") missing_columns = [col for col in required_columns if col not in reader.fieldnames] if missing_columns: raise ValueError( f"Missing required columns: {', '.join(missing_columns)}\n" f"Available columns: {', '.join(reader.fieldnames)}" ) for row_num, row in enumerate(reader, start=2): partial_id = row.get('partial_id', '').strip() partial_ordinal = row.get('partial_ordinal', '').strip() partial_last = row.get('partial_last', '').strip() content = row.get('Content', '') # If no partial_id, this is a complete log if not partial_id: complete_logs.append({ 'content': content, 'source': 'complete', 'row': row_num }) continue # Store fragment with metadata try: ordinal = int(partial_ordinal) if partial_ordinal else 0 except ValueError: ordinal = 0 log_fragments[partial_id].append({ 'ordinal': ordinal, 'content': content, 'is_last': partial_last.lower() == 'true', 'row': row_num }) # Reassemble fragments reassembled_logs = [] incomplete_logs = [] for partial_id, fragments in log_fragments.items(): # Sort by ordinal fragments.sort(key=lambda x: x['ordinal']) # Check if we have all parts last_fragment = None for frag in fragments: if frag['is_last']: last_fragment = frag break # Expected number of parts (ordinals start at 1) if last_fragment: expected_parts = last_fragment['ordinal'] actual_parts = len(fragments) # Check if we have all parts ordinals = [f['ordinal'] for f in fragments] expected_ordinals = list(range(1, expected_parts + 1)) if ordinals == expected_ordinals: # Reassemble the log reassembled_content = ''.join(f['content'] for f in fragments) # Try to parse as JSON to verify is_valid_json = False try: json.loads(reassembled_content) is_valid_json = True except json.JSONDecodeError: pass reassembled_logs.append({ 'partial_id': partial_id, 'content': reassembled_content, 'num_parts': actual_parts, 'is_valid_json': is_valid_json, 'fragments': fragments }) else: # Missing parts incomplete_logs.append({ 'partial_id': partial_id, 'expected_parts': expected_parts, 'actual_parts': actual_parts, 'missing_ordinals': sorted(set(expected_ordinals) - set(ordinals)), 'fragments': fragments }) else: # No last fragment found incomplete_logs.append({ 'partial_id': partial_id, 'expected_parts': 'unknown', 'actual_parts': len(fragments), 'reason': 'no_last_fragment', 'fragments': fragments }) # Combine complete and reassembled logs all_complete = complete_logs + [ { 'content': log['content'], 'source': 'reassembled', 'partial_id': log['partial_id'], 'num_parts': log['num_parts'], 'is_valid_json': log['is_valid_json'] } for log in reassembled_logs ] stats = { 'total_complete_logs': len(complete_logs), 'total_reassembled_logs': len(reassembled_logs), 'total_incomplete_logs': len(incomplete_logs), 'valid_json_reassembled': sum(1 for log in reassembled_logs if log['is_valid_json']), 'invalid_json_reassembled': sum(1 for log in reassembled_logs if not log['is_valid_json']) } return all_complete, incomplete_logs, stats def write_output(complete_logs, incomplete_logs, stats, output_dir): """Write reassembled logs and report to output directory.""" output_path = Path(output_dir) output_path.mkdir(parents=True, exist_ok=True) # Write complete logs (one per line, as JSONL) complete_file = output_path / 'reassembled_logs.jsonl' with open(complete_file, 'w', encoding='utf-8') as f: for log in complete_logs: # Try to parse and write as proper JSON try: parsed = json.loads(log['content']) f.write(json.dumps(parsed) + '\n') except json.JSONDecodeError: # Write as-is if not valid JSON f.write(log['content'] + '\n') # Write report report_file = output_path / 'reassembly_report.txt' with open(report_file, 'w', encoding='utf-8') as f: f.write("=" * 80 + "\n") f.write("Log Reassembly Report\n") f.write("=" * 80 + "\n\n") f.write(f"Complete logs (no splitting): {stats['total_complete_logs']}\n") f.write(f"Successfully reassembled logs: {stats['total_reassembled_logs']}\n") f.write(f" - Valid JSON: {stats['valid_json_reassembled']}\n") f.write(f" - Invalid JSON: {stats['invalid_json_reassembled']}\n") f.write(f"Incomplete logs (missing parts): {stats['total_incomplete_logs']}\n") f.write(f"\nTotal complete logs written: {len(complete_logs)}\n") if incomplete_logs: f.write("\n" + "=" * 80 + "\n") f.write("Incomplete Logs Details\n") f.write("=" * 80 + "\n\n") for i, log in enumerate(incomplete_logs, 1): f.write(f"Incomplete Log #{i}:\n") f.write(f" Partial ID: {log['partial_id']}\n") f.write(f" Expected parts: {log['expected_parts']}\n") f.write(f" Actual parts: {log['actual_parts']}\n") if 'missing_ordinals' in log: f.write(f" Missing ordinals: {log['missing_ordinals']}\n") if 'reason' in log: f.write(f" Reason: {log['reason']}\n") f.write(f" Available ordinals: {sorted([f['ordinal'] for f in log['fragments']])}\n") f.write("\n") return complete_file, report_file def main(): # Check for command-line argument if len(sys.argv) > 1: csv_file = Path(sys.argv[1]) if not csv_file.exists(): print(f"Error: File not found: {csv_file}") return else: # Find the most recent CSV file in the split_logs directory csv_files = sorted(Path('split_logs').glob('extract-*.csv'), reverse=True) if not csv_files: print("Error: No CSV file found matching pattern 'split_logs/extract-*.csv'") print("Usage: python3 reassemble_logs.py ") return csv_file = csv_files[0] # Determine output directory (same directory as input file) output_dir = csv_file.parent / 'output' print(f"Processing file: {csv_file}") print(f"Output directory: {output_dir}") print() # Reassemble logs try: complete_logs, incomplete_logs, stats = reassemble_split_logs(csv_file, output_dir) except ValueError as e: print(f"Error: {e}") return except Exception as e: print(f"Unexpected error: {e}") return # Write output complete_file, report_file = write_output(complete_logs, incomplete_logs, stats, output_dir) # Display summary print("=" * 80) print("Log Reassembly Complete") print("=" * 80) print(f"\nComplete logs (no splitting): {stats['total_complete_logs']}") print(f"Successfully reassembled logs: {stats['total_reassembled_logs']}") print(f" - Valid JSON: {stats['valid_json_reassembled']}") print(f" - Invalid JSON: {stats['invalid_json_reassembled']}") print(f"Incomplete logs (missing parts): {stats['total_incomplete_logs']}") print(f"\nTotal complete logs written: {len(complete_logs)}") print(f"\nOutput files:") print(f" - {complete_file}") print(f" - {report_file}") if stats['total_incomplete_logs'] > 0: print(f"\n⚠️ Warning: {stats['total_incomplete_logs']} logs could not be fully reassembled") print(f" See {report_file} for details") if __name__ == '__main__': main()