feat: implement legacy history migration (closes #24)
This commit is contained in:
@@ -0,0 +1,440 @@
|
||||
"""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 <path> (§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,
|
||||
}
|
||||
Reference in New Issue
Block a user