"""Legacy migration: import history.jsonl into the observation store. Spec §3.5, ADR 0001 §6. Idempotent and interruption-safe. Entry points: - Installer import detection at ./data/history.jsonl (§10.1) - fenris import (§8.8) - Collector's first new-version run (ADR 0001 §6) The import is a single transaction — a scripted kill mid-import leaves the store fully pre- or fully post-migration. history.jsonl is the sole authority; hourly.jsonl is diffed and logged but never trusted. Malformed lines are quarantined with a logged count, never silently dropped. Legacy files are renamed *.migrated only after commit and never deleted. """ import json import logging import sqlite3 from datetime import datetime, timezone from pathlib import Path from typing import Any, Dict, List, Optional, Tuple from .hour_classify import classify_hour from .segment import open_segment, normalize_identity from .monitoring_periods import close_period logger = logging.getLogger(__name__) # Legacy import marker — stored in a metadata table LEGACY_IMPORT_MARKER = "legacy_imported" def _ensure_metadata_table(conn: sqlite3.Connection) -> None: """Create the metadata table if it doesn't exist.""" conn.execute(""" CREATE TABLE IF NOT EXISTS store_metadata ( key TEXT PRIMARY KEY, value TEXT NOT NULL ) """) def is_legacy_imported(conn: sqlite3.Connection) -> bool: """Check if legacy history has already been imported. Spec §3.5.1: If the store already carries the legacy-import marker, do nothing. """ _ensure_metadata_table(conn) cursor = conn.execute( "SELECT value FROM store_metadata WHERE key = ?", (LEGACY_IMPORT_MARKER,), ) row = cursor.fetchone() return row is not None and row[0] == "true" def _parse_history_line(line: str, line_num: int) -> Optional[Dict[str, Any]]: """Parse a single line from history.jsonl. Returns None for malformed lines (quarantined, not silently dropped). """ line = line.strip() if not line: return None try: record = json.loads(line) except json.JSONDecodeError as e: logger.warning("Malformed JSON at line %d: %s", line_num, e) return None # Validate required fields required_fields = [ "timestamp", "model", "serial", "firmware_version", "data_units_written", "data_units_read", "percentage_used", "power_on_hours", "temperature", ] for field in required_fields: if field not in record: logger.warning("Missing field '%s' at line %d", field, line_num) return None return record def _record_to_sample(record: Dict[str, Any]) -> Dict[str, Any]: """Convert a legacy history.jsonl record to a sample dict.""" # Legacy records use different field names duw = record["data_units_written"] dur = record["data_units_read"] return { "ts": record["timestamp"], "device": record.get("device", "/dev/nvme0"), "subnqn": record.get("subsystem_nqn", ""), "sn": record["serial"], "mn": record["model"], "fr": record["firmware_version"], "capacity_bytes": record.get("capacity_bytes", 0), "percentage_used": record["percentage_used"], "available_spare": record.get("available_spare"), "media_errors": record.get("media_errors", 0), "power_on_hours": record["power_on_hours"], "power_cycles": record.get("power_cycles"), "unsafe_shutdowns": record.get("unsafe_shutdowns"), "temperature_c": record["temperature"], "data_units_written": duw, "data_units_read": dur, "bytes_written": duw * 512000, "bytes_read": dur * 512000, "critical_warning": record.get("critical_warning", 0), } def _derive_hour_observation( samples: List[Dict[str, Any]], hour_start: datetime, ) -> Dict[str, Any]: """Derive a single hour observation from samples in that hour. Pre-migration hours carry an unknown activity split except directly evidenced facts — a sample present means powered on; a DUW delta means writes occurred (§3.5.4). """ hour_end = hour_start.replace(hour=hour_start.hour + 1) if hour_start.hour < 23 else hour_start.replace(hour=0, day=hour_start.day + 1) # Filter samples in this hour hour_samples = [] for s in samples: ts = datetime.fromisoformat(s["ts"]) if hour_start <= ts < hour_end: hour_samples.append(s) if not hour_samples: return None # Sort by timestamp hour_samples.sort(key=lambda x: x["ts"]) # Compute deltas from first to last sample in the hour first = hour_samples[0] last = hour_samples[-1] duw_delta = last["bytes_written"] - first["bytes_written"] dur_delta = last["bytes_read"] - first["bytes_read"] poh_delta = (last["power_on_hours"] - first["power_on_hours"]) * 3600 # Temperature stats temps = [s["temperature_c"] for s in hour_samples] # Classify the hour split = classify_hour( wall_clock_seconds=3600, poh_delta=poh_delta, duw_delta=duw_delta, dur_delta=dur_delta, sampled_seconds=3600, # Legacy samples cover the full hour ) return { "hour": hour_start.strftime("%Y-%m-%dT%H:00:00Z"), "active_seconds": split.seconds_active, "idle_seconds": split.seconds_idle, "powered_off_seconds": split.seconds_powered_off, "unknown_seconds": split.seconds_unknown, "bytes_written_delta": duw_delta, "bytes_read_delta": dur_delta, "temperature_min": min(temps), "temperature_avg": sum(temps) / len(temps), "temperature_max": max(temps), "sample_count": len(hour_samples), "coverage": 1.0 if split.seconds_unknown == 0 else (3600 - split.seconds_unknown) / 3600, } def _diff_hourly_jsonl( hourly_path: Path, derived_hours: Dict[str, Dict[str, Any]], ) -> None: """Diff hourly.jsonl against derived data and log mismatches. Spec §3.5.3: hourly.jsonl is never trusted; mismatches are diffed and logged. """ if not hourly_path.exists(): logger.info("No hourly.jsonl found for diffing") return try: with open(hourly_path, "r") as f: for line_num, line in enumerate(f, 1): line = line.strip() if not line: continue try: record = json.loads(line) except json.JSONDecodeError: logger.warning("Malformed hourly.jsonl at line %d", line_num) continue hour_key = record.get("hour") if hour_key not in derived_hours: logger.info("Hourly.jsonl has hour %s not in derived data", hour_key) continue derived = derived_hours[hour_key] mismatches = [] for field in ["bytes_written_delta", "bytes_read_delta", "sample_count"]: if field in record and record[field] != derived.get(field): mismatches.append( f"{field}: hourly={record[field]} derived={derived.get(field)}" ) if mismatches: logger.info( "Hourly.jsonl mismatch for %s: %s", hour_key, "; ".join(mismatches), ) except Exception as e: logger.warning("Failed to diff hourly.jsonl: %s", e) def import_legacy_history( conn: sqlite3.Connection, history_path: Path, hourly_path: Optional[Path] = None, clock=None, ) -> Dict[str, Any]: """Import legacy history.jsonl into the observation store. This is the main entry point for legacy migration. It is: - Idempotent: second run no-ops on the legacy-import marker - Interruption-safe: single transaction - Never creates synthetic baselines Args: conn: Connection to the observation store history_path: Path to history.jsonl hourly_path: Optional path to hourly.jsonl for diffing clock: Injected clock (for testing) Returns: Dict with migration outcome """ # Check idempotency if is_legacy_imported(conn): return {"ok": True, "skipped": True, "reason": "already_imported"} # Read and parse history.jsonl samples = [] malformed_count = 0 if not history_path.exists(): return {"ok": False, "error": f"History file not found: {history_path}"} with open(history_path, "r") as f: for line_num, line in enumerate(f, 1): record = _parse_history_line(line, line_num) if record is None: malformed_count += 1 continue samples.append(_record_to_sample(record)) if malformed_count > 0: logger.warning("Quarantined %d malformed lines from history.jsonl", malformed_count) if not samples: return {"ok": False, "error": "No valid samples found in history.jsonl"} # Sort samples by timestamp samples.sort(key=lambda x: x["ts"]) # Derive hour observations hour_observations = {} for sample in samples: ts = datetime.fromisoformat(sample["ts"]) hour_start = ts.replace(minute=0, second=0, microsecond=0) hour_key = hour_start.strftime("%Y-%m-%dT%H:00:00Z") if hour_key not in hour_observations: hour_observations[hour_key] = { "hour": hour_start, "samples": [], } hour_observations[hour_key]["samples"].append(sample) derived_hours = {} for hour_key, hour_data in hour_observations.items(): obs = _derive_hour_observation(hour_data["samples"], hour_data["hour"]) if obs is not None: derived_hours[hour_key] = obs # Diff against hourly.jsonl if provided if hourly_path: _diff_hourly_jsonl(hourly_path, derived_hours) # Single transaction for the entire import try: # Begin transaction conn.execute("BEGIN IMMEDIATE") # 1. Insert raw samples for sample in samples: conn.execute( """ INSERT INTO samples ( ts, device, subnqn, sn, mn, fr, capacity_bytes, percentage_used, available_spare, media_errors, power_on_hours, power_cycles, unsafe_shutdowns, temperature_c, data_units_written, data_units_read, bytes_written, bytes_read, critical_warning ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) """, ( sample["ts"], sample["device"], sample["subnqn"], sample["sn"], sample["mn"], sample["fr"], sample["capacity_bytes"], sample["percentage_used"], sample["available_spare"], sample["media_errors"], sample["power_on_hours"], sample["power_cycles"], sample["unsafe_shutdowns"], sample["temperature_c"], sample["data_units_written"], sample["data_units_read"], sample["bytes_written"], sample["bytes_read"], sample["critical_warning"], ), ) # 2. Insert derived hour observations for hour_key, obs in sorted(derived_hours.items()): conn.execute( """ INSERT OR REPLACE INTO hour_observations ( hour, active_seconds, idle_seconds, powered_off_seconds, unknown_seconds, bytes_written_delta, bytes_read_delta, temperature_min, temperature_avg, temperature_max, sample_count, coverage ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) """, ( obs["hour"], obs["active_seconds"], obs["idle_seconds"], obs["powered_off_seconds"], obs["unknown_seconds"], obs["bytes_written_delta"], obs["bytes_read_delta"], obs["temperature_min"], obs["temperature_avg"], obs["temperature_max"], obs["sample_count"], obs["coverage"], ), ) # 3. Open implicit monitoring period at first legacy sample first_sample_ts = datetime.fromisoformat(samples[0]["ts"]) conn.execute( "INSERT INTO monitoring_periods (started_at) VALUES (?)", (first_sample_ts.isoformat(),), ) # 4. Close period with end_cause = migrated at migration moment if clock: migration_time = clock.utcnow() else: migration_time = datetime.now(timezone.utc) conn.execute( "UPDATE monitoring_periods SET ended_at = ?, end_cause = ? WHERE ended_at IS NULL", (migration_time.isoformat(), "migrated"), ) # 5. Open legacy controller segment (mn-only) # Legacy identity is model-scoped only (§4.4) first_sample = samples[0] legacy_identity = { "mn": first_sample["mn"], "sn": "", # Legacy segments are mn-only "subnqn": "", "fr": "", } legacy_identity_key = f"legacy|{first_sample['mn']}" open_segment( conn, migration_time, legacy_identity, legacy_identity_key, identity_degraded=False, ) # 6. Set legacy import marker _ensure_metadata_table(conn) conn.execute( "INSERT OR REPLACE INTO store_metadata (key, value) VALUES (?, ?)", (LEGACY_IMPORT_MARKER, "true"), ) # Commit conn.commit() except Exception as e: conn.rollback() raise RuntimeError(f"Migration failed: {e}") from e # 7. Rename legacy files to *.migrated (only after commit) try: migrated_path = history_path.with_suffix(history_path.suffix + ".migrated") history_path.rename(migrated_path) logger.info("Renamed %s to %s", history_path, migrated_path) if hourly_path and hourly_path.exists(): hourly_migrated = hourly_path.with_suffix(hourly_path.suffix + ".migrated") hourly_path.rename(hourly_migrated) logger.info("Renamed %s to %s", hourly_path, hourly_migrated) except Exception as e: # Non-fatal: files weren't renamed but migration succeeded logger.warning("Failed to rename legacy files: %s", e) return { "ok": True, "skipped": False, "samples_imported": len(samples), "hours_imported": len(derived_hours), "malformed_lines": malformed_count, "first_sample": samples[0]["ts"], "last_sample": samples[-1]["ts"], "legacy_identity_key": legacy_identity_key, }