Compare commits
3
Commits
7c21b044ba
...
8a5ef05188
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
8a5ef05188 | ||
|
|
847bde1be0 | ||
|
|
0d45f263d2 |
+9
-10
@@ -27,7 +27,6 @@ if VENV_DIR.exists():
|
|||||||
if site_packages:
|
if site_packages:
|
||||||
sys.path.insert(0, str(site_packages))
|
sys.path.insert(0, str(site_packages))
|
||||||
|
|
||||||
from fenris.store import init_store, get_store_path
|
|
||||||
from fenris.collector import run_collection
|
from fenris.collector import run_collection
|
||||||
|
|
||||||
|
|
||||||
@@ -103,14 +102,6 @@ def main() -> None:
|
|||||||
config = load_config()
|
config = load_config()
|
||||||
device = config["device"]
|
device = config["device"]
|
||||||
|
|
||||||
# Interrogate the drive
|
|
||||||
smartctl_data = interrogate_drive(device)
|
|
||||||
|
|
||||||
# Find sysfs path
|
|
||||||
sysfs_path = find_nvme_sysfs()
|
|
||||||
if sysfs_path is None:
|
|
||||||
raise RuntimeError("No NVMe controller found in sysfs")
|
|
||||||
|
|
||||||
# Inject a simple clock
|
# Inject a simple clock
|
||||||
class SimpleClock:
|
class SimpleClock:
|
||||||
def utcnow(self):
|
def utcnow(self):
|
||||||
@@ -118,8 +109,16 @@ def main() -> None:
|
|||||||
|
|
||||||
clock = SimpleClock()
|
clock = SimpleClock()
|
||||||
|
|
||||||
|
def acquire():
|
||||||
|
"""Interrogate the drive after pending-work preflight."""
|
||||||
|
smartctl_data = interrogate_drive(device)
|
||||||
|
sysfs_path = find_nvme_sysfs()
|
||||||
|
if sysfs_path is None:
|
||||||
|
raise RuntimeError("No NVMe controller found in sysfs")
|
||||||
|
return smartctl_data, sysfs_path
|
||||||
|
|
||||||
# Run collection
|
# Run collection
|
||||||
result = run_collection(smartctl_data, sysfs_path, config, clock)
|
result = run_collection(config=config, clock=clock, acquire=acquire)
|
||||||
|
|
||||||
if result["ok"]:
|
if result["ok"]:
|
||||||
print(f"Collection successful: {result['sample_count']} sample(s)")
|
print(f"Collection successful: {result['sample_count']} sample(s)")
|
||||||
|
|||||||
+177
-65
@@ -5,19 +5,21 @@ This module implements the thinnest complete write path:
|
|||||||
- Acquire controller identity from sysfs
|
- Acquire controller identity from sysfs
|
||||||
- Normalize identity exactly once at write time
|
- Normalize identity exactly once at write time
|
||||||
- Validate every row against store invariants
|
- Validate every row against store invariants
|
||||||
|
- Stage valid observations privately before derivation
|
||||||
- Publish the sample and derived evidence in one collection-owned transaction
|
- Publish the sample and derived evidence in one collection-owned transaction
|
||||||
|
|
||||||
No code path outside the collector interrogates the device.
|
No code path outside the collector interrogates the device.
|
||||||
"""
|
"""
|
||||||
import json
|
import json
|
||||||
import sqlite3
|
import sqlite3
|
||||||
|
from collections.abc import Callable
|
||||||
from datetime import datetime, timezone
|
from datetime import datetime, timezone
|
||||||
from pathlib import Path
|
from pathlib import Path
|
||||||
from typing import Any, Dict, Optional, Tuple
|
from typing import Any, Dict, Optional
|
||||||
|
|
||||||
from .store import init_store, get_store_path
|
from .derive import derive_hours_from_interval, find_previous_sample
|
||||||
from .monitoring_periods import ensure_period_open
|
from .monitoring_periods import ensure_period_open
|
||||||
from .derive import find_previous_sample, derive_hours_from_interval
|
from .store import init_store, get_store_path
|
||||||
|
|
||||||
|
|
||||||
class AcquisitionError(Exception):
|
class AcquisitionError(Exception):
|
||||||
@@ -30,6 +32,9 @@ class InvariantViolationError(Exception):
|
|||||||
pass
|
pass
|
||||||
|
|
||||||
|
|
||||||
|
PENDING_PUBLICATION_LIMIT = 6720
|
||||||
|
|
||||||
|
|
||||||
def acquire_from_smartctl(smartctl_data: Dict[str, Any]) -> Dict[str, Any]:
|
def acquire_from_smartctl(smartctl_data: Dict[str, Any]) -> Dict[str, Any]:
|
||||||
"""Acquire counters and thermal evidence from smartctl -a -j data.
|
"""Acquire counters and thermal evidence from smartctl -a -j data.
|
||||||
|
|
||||||
@@ -195,7 +200,7 @@ def write_sample(
|
|||||||
sample: Dict[str, Any],
|
sample: Dict[str, Any],
|
||||||
identity: Dict[str, Any],
|
identity: Dict[str, Any],
|
||||||
conn: sqlite3.Connection,
|
conn: sqlite3.Connection,
|
||||||
clock,
|
observed_at: datetime,
|
||||||
) -> Dict[str, Any]:
|
) -> Dict[str, Any]:
|
||||||
"""Write one sample to the observation store.
|
"""Write one sample to the observation store.
|
||||||
|
|
||||||
@@ -219,8 +224,7 @@ def write_sample(
|
|||||||
# Open new segment if needed
|
# Open new segment if needed
|
||||||
segment_opened = False
|
segment_opened = False
|
||||||
if should_open:
|
if should_open:
|
||||||
now = clock.utcnow()
|
open_segment(conn, observed_at, identity, identity_key, identity_degraded)
|
||||||
open_segment(conn, now, identity, identity_key, identity_degraded)
|
|
||||||
segment_opened = True
|
segment_opened = True
|
||||||
|
|
||||||
# Get current segment_id for provenance
|
# Get current segment_id for provenance
|
||||||
@@ -271,91 +275,199 @@ def write_sample(
|
|||||||
}
|
}
|
||||||
|
|
||||||
|
|
||||||
|
def _observation_time(sample: Dict[str, Any]) -> datetime:
|
||||||
|
"""Read an observation's original timestamp for ordered recovery."""
|
||||||
|
observed_at = datetime.fromisoformat(sample["ts"])
|
||||||
|
if observed_at.tzinfo is None:
|
||||||
|
return observed_at.replace(tzinfo=timezone.utc)
|
||||||
|
return observed_at.astimezone(timezone.utc)
|
||||||
|
|
||||||
|
|
||||||
|
def _publish_observation(
|
||||||
|
conn: sqlite3.Connection,
|
||||||
|
sample: Dict[str, Any],
|
||||||
|
identity: Dict[str, Any],
|
||||||
|
tz_name: str,
|
||||||
|
) -> None:
|
||||||
|
"""Publish one staged sample and all dependent evidence in caller transaction."""
|
||||||
|
observed_at = _observation_time(sample)
|
||||||
|
validate_sample_invariants(sample, conn)
|
||||||
|
ensure_period_open(conn, observed_at)
|
||||||
|
seg_info = write_sample(sample, identity, conn, observed_at)
|
||||||
|
|
||||||
|
cursor = conn.execute("SELECT id FROM samples ORDER BY id DESC LIMIT 1")
|
||||||
|
current_id = cursor.fetchone()[0]
|
||||||
|
prev = find_previous_sample(conn, seg_info.get("segment_id"), current_id)
|
||||||
|
if prev is not None:
|
||||||
|
current = {
|
||||||
|
"id": current_id,
|
||||||
|
"ts": sample["ts"],
|
||||||
|
"bytes_written": sample["bytes_written"],
|
||||||
|
"bytes_read": sample["bytes_read"],
|
||||||
|
"power_on_hours": sample["power_on_hours"],
|
||||||
|
"temperature_c": sample["temperature_c"],
|
||||||
|
"data_units_written": sample["data_units_written"],
|
||||||
|
"data_units_read": sample["data_units_read"],
|
||||||
|
}
|
||||||
|
derive_hours_from_interval(conn, prev, current)
|
||||||
|
|
||||||
|
from .day_aggregate import derive_all_days, persist_day_aggregate
|
||||||
|
for aggregate in derive_all_days(conn):
|
||||||
|
persist_day_aggregate(conn, aggregate)
|
||||||
|
|
||||||
|
from .local_day import derive_local_day_summary, persist_local_day
|
||||||
|
local_summary = derive_local_day_summary(conn, tz_name, observed_at)
|
||||||
|
if local_summary is not None:
|
||||||
|
persist_local_day(conn, local_summary)
|
||||||
|
|
||||||
|
|
||||||
|
def _pending_count(conn: sqlite3.Connection) -> int:
|
||||||
|
return conn.execute("SELECT COUNT(*) FROM pending_publications").fetchone()[0]
|
||||||
|
|
||||||
|
|
||||||
|
def _stage_observation(
|
||||||
|
conn: sqlite3.Connection,
|
||||||
|
sample: Dict[str, Any],
|
||||||
|
identity: Dict[str, Any],
|
||||||
|
tz_name: str,
|
||||||
|
) -> int:
|
||||||
|
"""Durably stage valid acquired evidence before attempting publication."""
|
||||||
|
payload = json.dumps(
|
||||||
|
{"sample": sample, "identity": identity, "tz_name": tz_name},
|
||||||
|
separators=(",", ":"),
|
||||||
|
sort_keys=True,
|
||||||
|
)
|
||||||
|
cursor = conn.execute(
|
||||||
|
"INSERT INTO pending_publications (sample_ts, payload) VALUES (?, ?)",
|
||||||
|
(sample["ts"], payload),
|
||||||
|
)
|
||||||
|
return cursor.lastrowid
|
||||||
|
|
||||||
|
|
||||||
|
def _recover_pending(conn: sqlite3.Connection) -> int:
|
||||||
|
"""Publish pending observations oldest first; stop at first failure."""
|
||||||
|
recovered = 0
|
||||||
|
while True:
|
||||||
|
conn.execute("BEGIN IMMEDIATE")
|
||||||
|
try:
|
||||||
|
row = conn.execute(
|
||||||
|
"SELECT id, payload FROM pending_publications ORDER BY id LIMIT 1"
|
||||||
|
).fetchone()
|
||||||
|
if row is None:
|
||||||
|
conn.rollback()
|
||||||
|
return recovered
|
||||||
|
|
||||||
|
pending_id, payload = row
|
||||||
|
observation = json.loads(payload)
|
||||||
|
_publish_observation(
|
||||||
|
conn,
|
||||||
|
observation["sample"],
|
||||||
|
observation["identity"],
|
||||||
|
observation["tz_name"],
|
||||||
|
)
|
||||||
|
conn.execute("DELETE FROM pending_publications WHERE id = ?", (pending_id,))
|
||||||
|
conn.commit()
|
||||||
|
recovered += 1
|
||||||
|
except Exception:
|
||||||
|
conn.rollback()
|
||||||
|
raise
|
||||||
|
|
||||||
|
|
||||||
|
def _recover_pending_for_collection(conn: sqlite3.Connection) -> int:
|
||||||
|
"""Retry queued work and report exhausted capacity without masking store faults."""
|
||||||
|
try:
|
||||||
|
return _recover_pending(conn)
|
||||||
|
except Exception as exc:
|
||||||
|
if isinstance(exc, sqlite3.Error) and not isinstance(exc, sqlite3.IntegrityError):
|
||||||
|
raise
|
||||||
|
if isinstance(exc, InvariantViolationError):
|
||||||
|
raise
|
||||||
|
try:
|
||||||
|
capacity_full = _pending_count(conn) >= PENDING_PUBLICATION_LIMIT
|
||||||
|
except sqlite3.Error:
|
||||||
|
raise exc
|
||||||
|
if capacity_full:
|
||||||
|
raise RuntimeError(
|
||||||
|
"pending publication capacity full "
|
||||||
|
f"({PENDING_PUBLICATION_LIMIT} observations); no new observation acquired; "
|
||||||
|
f"recovery failed: {exc}"
|
||||||
|
) from exc
|
||||||
|
raise
|
||||||
|
|
||||||
|
|
||||||
def run_collection(
|
def run_collection(
|
||||||
smartctl_data: Dict[str, Any],
|
smartctl_data: Optional[Dict[str, Any]] = None,
|
||||||
sysfs_path: Path,
|
sysfs_path: Optional[Path] = None,
|
||||||
config: Dict[str, Any],
|
config: Optional[Dict[str, Any]] = None,
|
||||||
clock,
|
clock=None,
|
||||||
|
*,
|
||||||
|
acquire: Optional[Callable[[], tuple[Dict[str, Any], Path]]] = None,
|
||||||
) -> Dict[str, Any]:
|
) -> Dict[str, Any]:
|
||||||
"""Run one collection run.
|
"""Recover old work, acquire one observation, then publish it atomically.
|
||||||
|
|
||||||
This is the main entry point for the collector.
|
Production callers pass ``acquire`` so recovery and capacity checks run
|
||||||
Returns the run outcome.
|
before device interrogation. Direct sample arguments remain useful for
|
||||||
|
deterministic collector tests.
|
||||||
"""
|
"""
|
||||||
conn = None
|
conn = None
|
||||||
try:
|
try:
|
||||||
# Acquire counters and thermal evidence
|
if config is None or clock is None:
|
||||||
counters = acquire_from_smartctl(smartctl_data)
|
raise ValueError("config and clock are required")
|
||||||
|
|
||||||
# Acquire controller identity
|
|
||||||
identity = acquire_from_sysfs(sysfs_path)
|
|
||||||
|
|
||||||
# Build sample with injected clock
|
|
||||||
sample = {
|
|
||||||
"ts": clock.utcnow().isoformat(),
|
|
||||||
"device": config["device"],
|
|
||||||
**counters,
|
|
||||||
**identity,
|
|
||||||
}
|
|
||||||
|
|
||||||
# Initialize store if needed
|
|
||||||
store_path = get_store_path(config)
|
store_path = get_store_path(config)
|
||||||
conn = init_store(store_path)
|
conn = init_store(store_path)
|
||||||
|
|
||||||
# Run legacy import if needed (idempotent)
|
|
||||||
from .legacy import import_legacy_history
|
from .legacy import import_legacy_history
|
||||||
history_path = Path(config.get("data_dir", ".")) / "history.jsonl"
|
history_path = Path(config.get("data_dir", ".")) / "history.jsonl"
|
||||||
if history_path.exists():
|
if history_path.exists():
|
||||||
import_legacy_history(conn, history_path, clock=clock)
|
import_legacy_history(conn, history_path, clock=clock)
|
||||||
|
|
||||||
# Collection owns one transaction for the sample and its evidence.
|
published_count = _recover_pending_for_collection(conn)
|
||||||
ensure_period_open(conn, clock.utcnow())
|
|
||||||
|
|
||||||
validate_sample_invariants(sample, conn)
|
# Hold the writer reservation across the capacity check and acquisition.
|
||||||
seg_info = write_sample(sample, identity, conn, clock)
|
# A concurrent collector will recheck pending work before it acquires.
|
||||||
|
while True:
|
||||||
|
conn.execute("BEGIN IMMEDIATE")
|
||||||
|
waiting = _pending_count(conn)
|
||||||
|
if waiting:
|
||||||
|
conn.rollback()
|
||||||
|
published_count += _recover_pending_for_collection(conn)
|
||||||
|
continue
|
||||||
|
|
||||||
cursor = conn.execute("SELECT id FROM samples ORDER BY id DESC LIMIT 1")
|
from .tz_util import detect_system_tz
|
||||||
current_id = cursor.fetchone()[0]
|
tz_name = detect_system_tz()
|
||||||
prev = find_previous_sample(conn, seg_info.get("segment_id"), current_id)
|
if acquire is not None:
|
||||||
if prev is not None:
|
smartctl_data, sysfs_path = acquire()
|
||||||
current = {
|
if smartctl_data is None or sysfs_path is None:
|
||||||
"id": current_id,
|
raise AcquisitionError("No acquired SMART data or sysfs identity")
|
||||||
"ts": sample["ts"],
|
|
||||||
"bytes_written": sample["bytes_written"],
|
counters = acquire_from_smartctl(smartctl_data)
|
||||||
"bytes_read": sample["bytes_read"],
|
identity = acquire_from_sysfs(sysfs_path)
|
||||||
"power_on_hours": sample["power_on_hours"],
|
sample = {
|
||||||
"temperature_c": sample["temperature_c"],
|
"ts": clock.utcnow().isoformat(),
|
||||||
"data_units_written": sample["data_units_written"],
|
"device": config["device"],
|
||||||
"data_units_read": sample["data_units_read"],
|
**counters,
|
||||||
|
**identity,
|
||||||
}
|
}
|
||||||
derive_hours_from_interval(conn, prev, current)
|
validate_sample_invariants(sample, conn)
|
||||||
|
_stage_observation(conn, sample, identity, tz_name)
|
||||||
from .day_aggregate import derive_all_days, persist_day_aggregate
|
conn.commit()
|
||||||
for agg in derive_all_days(conn):
|
break
|
||||||
persist_day_aggregate(conn, agg)
|
|
||||||
|
|
||||||
from .tz_util import detect_system_tz
|
|
||||||
from .local_day import derive_local_day_summary, persist_local_day
|
|
||||||
tz_name = detect_system_tz()
|
|
||||||
local_summary = derive_local_day_summary(conn, tz_name, clock.utcnow())
|
|
||||||
if local_summary is not None:
|
|
||||||
persist_local_day(conn, local_summary)
|
|
||||||
|
|
||||||
conn.commit()
|
|
||||||
|
|
||||||
|
published_count += _recover_pending_for_collection(conn)
|
||||||
return {
|
return {
|
||||||
"ok": True,
|
"ok": True,
|
||||||
"sample_count": 1,
|
"sample_count": published_count,
|
||||||
"store_path": str(store_path),
|
"store_path": str(store_path),
|
||||||
}
|
}
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as exc:
|
||||||
if conn is not None:
|
if conn is not None:
|
||||||
conn.rollback()
|
conn.rollback()
|
||||||
return {
|
return {
|
||||||
"ok": False,
|
"ok": False,
|
||||||
"error": str(e),
|
"error": str(exc),
|
||||||
"error_type": type(e).__name__,
|
"error_type": type(exc).__name__,
|
||||||
}
|
}
|
||||||
finally:
|
finally:
|
||||||
if conn is not None:
|
if conn is not None:
|
||||||
|
|||||||
@@ -115,6 +115,7 @@ class StatusComposition:
|
|||||||
# Sample counts for waiting explanations
|
# Sample counts for waiting explanations
|
||||||
sample_count: int = 0
|
sample_count: int = 0
|
||||||
day_count: int = 0
|
day_count: int = 0
|
||||||
|
pending_publication_count: int = 0
|
||||||
|
|
||||||
# Whether the status dot should blink (only Monitoring)
|
# Whether the status dot should blink (only Monitoring)
|
||||||
should_blink: bool = False
|
should_blink: bool = False
|
||||||
@@ -286,6 +287,7 @@ def compose_status(
|
|||||||
freshness_age_s = None
|
freshness_age_s = None
|
||||||
sample_count = 0
|
sample_count = 0
|
||||||
day_count = 0
|
day_count = 0
|
||||||
|
pending_publication_count = 0
|
||||||
deliberately_paused = False
|
deliberately_paused = False
|
||||||
|
|
||||||
if conn is None and store_fault is None and newer_schema is None:
|
if conn is None and store_fault is None and newer_schema is None:
|
||||||
@@ -315,6 +317,10 @@ def compose_status(
|
|||||||
sample_count = cursor.fetchone()[0]
|
sample_count = cursor.fetchone()[0]
|
||||||
cursor = conn.execute("SELECT COUNT(*) FROM day_aggregates")
|
cursor = conn.execute("SELECT COUNT(*) FROM day_aggregates")
|
||||||
day_count = cursor.fetchone()[0]
|
day_count = cursor.fetchone()[0]
|
||||||
|
schema_version = conn.execute("PRAGMA user_version").fetchone()[0]
|
||||||
|
if schema_version >= 4:
|
||||||
|
cursor = conn.execute("SELECT COUNT(*) FROM pending_publications")
|
||||||
|
pending_publication_count = cursor.fetchone()[0]
|
||||||
except sqlite3.Error as exc:
|
except sqlite3.Error as exc:
|
||||||
store_fault = str(exc)
|
store_fault = str(exc)
|
||||||
|
|
||||||
@@ -328,6 +334,7 @@ def compose_status(
|
|||||||
freshness = "unknown"
|
freshness = "unknown"
|
||||||
freshness_age_s = None
|
freshness_age_s = None
|
||||||
sample_count = day_count = 0
|
sample_count = day_count = 0
|
||||||
|
pending_publication_count = 0
|
||||||
deliberately_paused = False
|
deliberately_paused = False
|
||||||
|
|
||||||
# --- External stop detection ---
|
# --- External stop detection ---
|
||||||
@@ -417,6 +424,7 @@ def compose_status(
|
|||||||
overlay_base=overlay_base,
|
overlay_base=overlay_base,
|
||||||
sample_count=sample_count,
|
sample_count=sample_count,
|
||||||
day_count=day_count,
|
day_count=day_count,
|
||||||
|
pending_publication_count=pending_publication_count,
|
||||||
continuity=continuity,
|
continuity=continuity,
|
||||||
paused_lines=paused_lines,
|
paused_lines=paused_lines,
|
||||||
)
|
)
|
||||||
@@ -459,6 +467,8 @@ def render_status_cli(comp: StatusComposition) -> str:
|
|||||||
facts.append("timer: %s" % (
|
facts.append("timer: %s" % (
|
||||||
"unknown" if comp.timer_active is None else "active" if comp.timer_active else "inactive"
|
"unknown" if comp.timer_active is None else "active" if comp.timer_active else "inactive"
|
||||||
))
|
))
|
||||||
|
if comp.pending_publication_count:
|
||||||
|
facts.append(_pending_publication_text(comp.pending_publication_count))
|
||||||
|
|
||||||
if facts:
|
if facts:
|
||||||
lines.append(" · ".join(facts))
|
lines.append(" · ".join(facts))
|
||||||
@@ -524,6 +534,8 @@ def render_status_tui(comp: StatusComposition) -> str:
|
|||||||
facts.append("Timer: %s" % (
|
facts.append("Timer: %s" % (
|
||||||
"unknown" if comp.timer_active is None else "active" if comp.timer_active else "inactive"
|
"unknown" if comp.timer_active is None else "active" if comp.timer_active else "inactive"
|
||||||
))
|
))
|
||||||
|
if comp.pending_publication_count:
|
||||||
|
facts.append(_pending_publication_text(comp.pending_publication_count).capitalize())
|
||||||
|
|
||||||
if facts:
|
if facts:
|
||||||
lines.append(" · ".join(facts))
|
lines.append(" · ".join(facts))
|
||||||
@@ -539,3 +551,8 @@ def render_status_tui(comp: StatusComposition) -> str:
|
|||||||
lines.append(pl[:1].upper() + pl[1:])
|
lines.append(pl[:1].upper() + pl[1:])
|
||||||
|
|
||||||
return "\n".join(lines)
|
return "\n".join(lines)
|
||||||
|
|
||||||
|
|
||||||
|
def _pending_publication_text(count: int) -> str:
|
||||||
|
noun = "observation" if count == 1 else "observations"
|
||||||
|
return f"pending publication: {count} {noun} retained for retry"
|
||||||
|
|||||||
+22
-4
@@ -3,8 +3,7 @@
|
|||||||
This module handles:
|
This module handles:
|
||||||
- Store initialization with WAL mode
|
- Store initialization with WAL mode
|
||||||
- Schema versioning with PRAGMA user_version
|
- Schema versioning with PRAGMA user_version
|
||||||
- The six entities: samples, hour_observations, day_aggregates,
|
- Observation records, derived activity, publication state, and metadata
|
||||||
monitoring_periods, controller_segments, endurance_baseline
|
|
||||||
"""
|
"""
|
||||||
import sqlite3
|
import sqlite3
|
||||||
from pathlib import Path
|
from pathlib import Path
|
||||||
@@ -12,7 +11,7 @@ from typing import Optional
|
|||||||
|
|
||||||
|
|
||||||
# Schema version - increment on each migration
|
# Schema version - increment on each migration
|
||||||
SCHEMA_VERSION = 3
|
SCHEMA_VERSION = 4
|
||||||
|
|
||||||
|
|
||||||
# Packaged default placement (spec §8.3). The config may override it, but a
|
# Packaged default placement (spec §8.3). The config may override it, but a
|
||||||
@@ -33,7 +32,7 @@ def get_store_path(config: dict) -> Path:
|
|||||||
def init_store(store_path: Path) -> sqlite3.Connection:
|
def init_store(store_path: Path) -> sqlite3.Connection:
|
||||||
"""Initialize the observation store if not present.
|
"""Initialize the observation store if not present.
|
||||||
|
|
||||||
Creates the database with WAL mode and all six entities.
|
Creates the observation schema, including private pending-publication storage.
|
||||||
Returns a connection to the store.
|
Returns a connection to the store.
|
||||||
"""
|
"""
|
||||||
conn = sqlite3.connect(str(store_path))
|
conn = sqlite3.connect(str(store_path))
|
||||||
@@ -221,6 +220,19 @@ def _create_schema(conn: sqlite3.Connection):
|
|||||||
)
|
)
|
||||||
""")
|
""")
|
||||||
|
|
||||||
|
_create_pending_publications(conn)
|
||||||
|
|
||||||
|
|
||||||
|
def _create_pending_publications(conn: sqlite3.Connection) -> None:
|
||||||
|
"""Create private staging for valid observations awaiting derivation."""
|
||||||
|
conn.execute("""
|
||||||
|
CREATE TABLE IF NOT EXISTS pending_publications (
|
||||||
|
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||||||
|
sample_ts TEXT NOT NULL,
|
||||||
|
payload TEXT NOT NULL
|
||||||
|
)
|
||||||
|
""")
|
||||||
|
|
||||||
|
|
||||||
def _apply_migrations(conn: sqlite3.Connection, current_version: int):
|
def _apply_migrations(conn: sqlite3.Connection, current_version: int):
|
||||||
"""Apply forward-only migrations from current_version to SCHEMA_VERSION.
|
"""Apply forward-only migrations from current_version to SCHEMA_VERSION.
|
||||||
@@ -275,6 +287,12 @@ def _apply_migrations(conn: sqlite3.Connection, current_version: int):
|
|||||||
""")
|
""")
|
||||||
current_version = 3
|
current_version = 3
|
||||||
|
|
||||||
|
# Migration 3→4: retain acquired observations until derived evidence can
|
||||||
|
# be published atomically (issue #97, ADR 0011).
|
||||||
|
if current_version < 4:
|
||||||
|
_create_pending_publications(conn)
|
||||||
|
current_version = 4
|
||||||
|
|
||||||
|
|
||||||
def migrate_to_latest(store_path: Path) -> int:
|
def migrate_to_latest(store_path: Path) -> int:
|
||||||
"""Apply forward-only migrations to bring the store to SCHEMA_VERSION.
|
"""Apply forward-only migrations to bring the store to SCHEMA_VERSION.
|
||||||
|
|||||||
@@ -154,7 +154,7 @@ class TestSchemaMigration:
|
|||||||
|
|
||||||
# Migrate
|
# Migrate
|
||||||
steps = migrate_to_latest(db)
|
steps = migrate_to_latest(db)
|
||||||
assert steps == 2 # v1→v2→v3
|
assert steps == 3 # v1→v2→v3→v4
|
||||||
|
|
||||||
# Verify data preserved
|
# Verify data preserved
|
||||||
conn = sqlite3.connect(str(db))
|
conn = sqlite3.connect(str(db))
|
||||||
|
|||||||
@@ -232,7 +232,7 @@ def test_store_initialization(config_fixture: Dict[str, Any]):
|
|||||||
# Initialize store
|
# Initialize store
|
||||||
conn = init_store(store_path)
|
conn = init_store(store_path)
|
||||||
|
|
||||||
# Verify all six entities exist
|
# Verify the observation, derived-history, and metadata tables exist.
|
||||||
cursor = conn.execute("SELECT name FROM sqlite_master WHERE type='table'")
|
cursor = conn.execute("SELECT name FROM sqlite_master WHERE type='table'")
|
||||||
tables = {row[0] for row in cursor.fetchall()}
|
tables = {row[0] for row in cursor.fetchall()}
|
||||||
|
|
||||||
@@ -249,6 +249,7 @@ def test_store_initialization(config_fixture: Dict[str, Any]):
|
|||||||
expected_tables.add("sqlite_sequence")
|
expected_tables.add("sqlite_sequence")
|
||||||
expected_tables.add("store_metadata")
|
expected_tables.add("store_metadata")
|
||||||
expected_tables.add("local_days")
|
expected_tables.add("local_days")
|
||||||
|
expected_tables.add("pending_publications")
|
||||||
assert expected_tables == tables
|
assert expected_tables == tables
|
||||||
|
|
||||||
conn.close()
|
conn.close()
|
||||||
|
|||||||
@@ -0,0 +1,291 @@
|
|||||||
|
"""Pending publication behavior from issue #97."""
|
||||||
|
import sqlite3
|
||||||
|
import sys
|
||||||
|
from datetime import datetime, timedelta, timezone
|
||||||
|
from pathlib import Path
|
||||||
|
|
||||||
|
import pytest
|
||||||
|
|
||||||
|
sys.path.insert(0, str(Path(__file__).parent.parent / "src"))
|
||||||
|
|
||||||
|
from fenris import collector
|
||||||
|
from fenris.projection import compute_projection
|
||||||
|
from fenris.status import _query_drive_facts, get_status, read_status
|
||||||
|
from fenris.status_composition import render_status_tui
|
||||||
|
from fenris.store import init_store
|
||||||
|
|
||||||
|
|
||||||
|
class FakeClock:
|
||||||
|
def __init__(self, now):
|
||||||
|
self.now = now
|
||||||
|
|
||||||
|
def utcnow(self):
|
||||||
|
return self.now
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.fixture
|
||||||
|
def sysfs_fixture_tree(tmp_path):
|
||||||
|
ctrl_dir = tmp_path / "sys" / "class" / "nvme" / "nvme0"
|
||||||
|
ctrl_dir.mkdir(parents=True)
|
||||||
|
(ctrl_dir / "subsysnqn").write_text("nqn.test:drive\n")
|
||||||
|
(ctrl_dir / "model").write_text("Test NVMe\n")
|
||||||
|
(ctrl_dir / "serial").write_text("test-serial\n")
|
||||||
|
(ctrl_dir / "firmware_rev").write_text("1.0\n")
|
||||||
|
transport_dir = ctrl_dir / "transport"
|
||||||
|
transport_dir.mkdir()
|
||||||
|
(transport_dir / "trstring").write_text("pcie\n")
|
||||||
|
return tmp_path
|
||||||
|
|
||||||
|
|
||||||
|
def _smartctl(unit_written, unit_read=500):
|
||||||
|
return {
|
||||||
|
"nvme_smart_health_information_log": {
|
||||||
|
"critical_warning": 0,
|
||||||
|
"temperature": 35,
|
||||||
|
"available_spare": 100,
|
||||||
|
"percentage_used": 5,
|
||||||
|
"data_units_written": unit_written,
|
||||||
|
"data_units_read": unit_read,
|
||||||
|
"power_on_hours": 100,
|
||||||
|
},
|
||||||
|
"user_capacity": {"bytes": 1_024_000_000_000},
|
||||||
|
"model_name": "Test NVMe",
|
||||||
|
"serial_number": "test-serial",
|
||||||
|
"firmware_version": "1.0",
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
|
def _acquirer(sysfs_path, unit_written, unit_read=500):
|
||||||
|
return lambda: (_smartctl(unit_written, unit_read), sysfs_path)
|
||||||
|
|
||||||
|
|
||||||
|
def _public_snapshot(store_path, clock_now):
|
||||||
|
with read_status(store_path, clock_now) as (conn, comp):
|
||||||
|
assert conn is not None
|
||||||
|
return {
|
||||||
|
"state": comp.state,
|
||||||
|
"freshness": comp.freshness,
|
||||||
|
"sample_count": comp.sample_count,
|
||||||
|
"day_count": comp.day_count,
|
||||||
|
"boot_enabled": comp.boot_enabled,
|
||||||
|
"timer_active": comp.timer_active,
|
||||||
|
"last_collect_ok": comp.last_collect_ok,
|
||||||
|
"latest_sample": tuple(conn.execute(
|
||||||
|
"SELECT ts, subnqn, mn, sn, fr, capacity_bytes, temperature_c, "
|
||||||
|
"available_spare, data_units_written, bytes_written, critical_warning, "
|
||||||
|
"media_errors, unsafe_shutdowns "
|
||||||
|
"FROM samples ORDER BY id DESC LIMIT 1"
|
||||||
|
).fetchone()),
|
||||||
|
"drive_facts": tuple(_query_drive_facts(conn)),
|
||||||
|
"projection": repr(compute_projection(conn, clock_now)),
|
||||||
|
"segments": conn.execute("SELECT COUNT(*) FROM controller_segments").fetchone()[0],
|
||||||
|
"periods": conn.execute("SELECT COUNT(*) FROM monitoring_periods").fetchone()[0],
|
||||||
|
"hours": tuple(conn.execute(
|
||||||
|
"SELECT hour, bytes_written_delta, bytes_read_delta, sample_count "
|
||||||
|
"FROM hour_observations ORDER BY hour"
|
||||||
|
).fetchall()),
|
||||||
|
"days": tuple(conn.execute(
|
||||||
|
"SELECT day, bytes_written_delta, bytes_read_delta "
|
||||||
|
"FROM day_aggregates ORDER BY day"
|
||||||
|
).fetchall()),
|
||||||
|
"local_days": tuple(conn.execute(
|
||||||
|
"SELECT local_date, bytes_written, bytes_read, sample_count "
|
||||||
|
"FROM local_days ORDER BY local_date, tz_name"
|
||||||
|
).fetchall()),
|
||||||
|
"pending": comp.pending_publication_count,
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
|
def test_failed_publication_survives_restart_and_readers_keep_last_consistent_view(
|
||||||
|
tmp_path, sysfs_fixture_tree, monkeypatch
|
||||||
|
):
|
||||||
|
"""A failed derivation stays private, visible as pending, and replays in order."""
|
||||||
|
monkeypatch.setenv("TZ", "UTC")
|
||||||
|
service = {"boot_enabled": True, "timer_active": True, "last_collect_ok": True}
|
||||||
|
monkeypatch.setattr("fenris.status.query_service_state", lambda: service)
|
||||||
|
now = datetime.now(timezone.utc).replace(second=0, microsecond=0)
|
||||||
|
store_path = tmp_path / "observations.db"
|
||||||
|
config = {"device": "/dev/nvme0", "store_path": str(store_path)}
|
||||||
|
sysfs_path = sysfs_fixture_tree / "sys" / "class" / "nvme" / "nvme0"
|
||||||
|
|
||||||
|
for units, at in ((1000, now - timedelta(minutes=10)), (1010, now - timedelta(minutes=5))):
|
||||||
|
result = collector.run_collection(
|
||||||
|
config=config,
|
||||||
|
clock=FakeClock(at),
|
||||||
|
acquire=_acquirer(sysfs_path, units),
|
||||||
|
)
|
||||||
|
assert result["ok"] is True, result
|
||||||
|
|
||||||
|
previous = _public_snapshot(store_path, now)
|
||||||
|
baseline_cli = get_status(store_path, now, query_services=False, query_journal=False)
|
||||||
|
original_derive = collector.derive_hours_from_interval
|
||||||
|
reader_during_write = {}
|
||||||
|
|
||||||
|
def fail_derivation(conn, previous_sample, next_sample):
|
||||||
|
# Read while writer has inserted the unpublished sample but not committed.
|
||||||
|
with read_status(store_path, now, query_services=False) as (reader, comp):
|
||||||
|
assert reader is not None
|
||||||
|
reader_during_write["samples"] = reader.execute(
|
||||||
|
"SELECT COUNT(*) FROM samples"
|
||||||
|
).fetchone()[0]
|
||||||
|
reader_during_write["pending"] = comp.pending_publication_count
|
||||||
|
reader_during_write["latest"] = tuple(reader.execute(
|
||||||
|
"SELECT data_units_written FROM samples ORDER BY id DESC LIMIT 1"
|
||||||
|
).fetchone())
|
||||||
|
raise RuntimeError("injected non-invariant derivation failure")
|
||||||
|
|
||||||
|
monkeypatch.setattr(collector, "derive_hours_from_interval", fail_derivation)
|
||||||
|
failed = collector.run_collection(
|
||||||
|
config=config,
|
||||||
|
clock=FakeClock(now),
|
||||||
|
acquire=_acquirer(sysfs_path, 1025),
|
||||||
|
)
|
||||||
|
assert failed["ok"] is False
|
||||||
|
assert "injected non-invariant derivation failure" in failed["error"]
|
||||||
|
assert reader_during_write == {"samples": 2, "pending": 1, "latest": (1010,)}
|
||||||
|
|
||||||
|
after_failure = _public_snapshot(store_path, now)
|
||||||
|
assert {**after_failure, "pending": previous["pending"]} == previous
|
||||||
|
assert after_failure["pending"] == 1
|
||||||
|
|
||||||
|
cli_text = get_status(store_path, now, query_services=False, query_journal=False)
|
||||||
|
pending_fact = " · pending publication: 1 observation retained for retry"
|
||||||
|
assert pending_fact in cli_text
|
||||||
|
assert cli_text.replace(pending_fact, "", 1) == baseline_cli
|
||||||
|
|
||||||
|
with read_status(store_path, now) as (_, comp):
|
||||||
|
assert comp.state == previous["state"]
|
||||||
|
assert comp.boot_enabled is True
|
||||||
|
assert comp.timer_active is True
|
||||||
|
assert "pending publication: 1 observation retained for retry" in render_status_tui(comp).lower()
|
||||||
|
|
||||||
|
called = False
|
||||||
|
|
||||||
|
def should_wait_for_recovery():
|
||||||
|
nonlocal called
|
||||||
|
called = True
|
||||||
|
return _smartctl(1030), sysfs_path
|
||||||
|
|
||||||
|
blocked = collector.run_collection(
|
||||||
|
config=config,
|
||||||
|
clock=FakeClock(now + timedelta(minutes=5)),
|
||||||
|
acquire=should_wait_for_recovery,
|
||||||
|
)
|
||||||
|
assert blocked["ok"] is False
|
||||||
|
assert called is False
|
||||||
|
assert _public_snapshot(store_path, now)["pending"] == 1
|
||||||
|
|
||||||
|
monkeypatch.setattr(collector, "derive_hours_from_interval", original_derive)
|
||||||
|
monkeypatch.setenv("TZ", "Asia/Kolkata")
|
||||||
|
recovered = collector.run_collection(
|
||||||
|
config=config,
|
||||||
|
clock=FakeClock(now + timedelta(minutes=5)),
|
||||||
|
acquire=_acquirer(sysfs_path, 1040),
|
||||||
|
)
|
||||||
|
assert recovered["ok"] is True, recovered
|
||||||
|
|
||||||
|
conn = sqlite3.connect(store_path)
|
||||||
|
try:
|
||||||
|
assert conn.execute(
|
||||||
|
"SELECT data_units_written FROM samples ORDER BY id"
|
||||||
|
).fetchall() == [(1000,), (1010,), (1025,), (1040,)]
|
||||||
|
assert conn.execute("SELECT COUNT(*) FROM pending_publications").fetchone()[0] == 0
|
||||||
|
assert conn.execute(
|
||||||
|
"SELECT SUM(bytes_written_delta) FROM day_aggregates"
|
||||||
|
).fetchone()[0] == (1040 - 1000) * 512_000
|
||||||
|
assert conn.execute(
|
||||||
|
"SELECT sample_count FROM local_days WHERE tz_name = 'UTC'"
|
||||||
|
).fetchone()[0] == 4
|
||||||
|
assert conn.execute(
|
||||||
|
"SELECT COUNT(*) FROM local_days WHERE tz_name = 'Asia/Kolkata'"
|
||||||
|
).fetchone()[0] == 1
|
||||||
|
finally:
|
||||||
|
conn.close()
|
||||||
|
|
||||||
|
|
||||||
|
def test_invalid_sample_writes_no_sample_segment_period_or_pending_row(tmp_path, sysfs_fixture_tree):
|
||||||
|
"""Sample invariants run before staging or public monitoring changes."""
|
||||||
|
now = datetime.now(timezone.utc).replace(second=0, microsecond=0)
|
||||||
|
store_path = tmp_path / "observations.db"
|
||||||
|
config = {"device": "/dev/nvme0", "store_path": str(store_path)}
|
||||||
|
sysfs_path = sysfs_fixture_tree / "sys" / "class" / "nvme" / "nvme0"
|
||||||
|
result = collector.run_collection(
|
||||||
|
config=config,
|
||||||
|
clock=FakeClock(now),
|
||||||
|
acquire=_acquirer(sysfs_path, -1),
|
||||||
|
)
|
||||||
|
|
||||||
|
assert result["ok"] is False
|
||||||
|
assert result["error_type"] == "InvariantViolationError"
|
||||||
|
conn = sqlite3.connect(store_path)
|
||||||
|
try:
|
||||||
|
assert conn.execute("SELECT COUNT(*) FROM samples").fetchone()[0] == 0
|
||||||
|
assert conn.execute("SELECT COUNT(*) FROM controller_segments").fetchone()[0] == 0
|
||||||
|
assert conn.execute("SELECT COUNT(*) FROM monitoring_periods").fetchone()[0] == 0
|
||||||
|
assert conn.execute("SELECT COUNT(*) FROM pending_publications").fetchone()[0] == 0
|
||||||
|
finally:
|
||||||
|
conn.close()
|
||||||
|
|
||||||
|
|
||||||
|
def test_pending_capacity_retries_before_skipping_acquisition(tmp_path, sysfs_fixture_tree, monkeypatch):
|
||||||
|
now = datetime.now(timezone.utc).replace(second=0, microsecond=0)
|
||||||
|
store_path = tmp_path / "observations.db"
|
||||||
|
config = {"device": "/dev/nvme0", "store_path": str(store_path)}
|
||||||
|
sysfs_path = sysfs_fixture_tree / "sys" / "class" / "nvme" / "nvme0"
|
||||||
|
first = collector.run_collection(
|
||||||
|
config=config,
|
||||||
|
clock=FakeClock(now - timedelta(minutes=5)),
|
||||||
|
acquire=_acquirer(sysfs_path, 1000),
|
||||||
|
)
|
||||||
|
assert first["ok"] is True, first
|
||||||
|
|
||||||
|
monkeypatch.setattr(collector, "PENDING_PUBLICATION_LIMIT", 1)
|
||||||
|
|
||||||
|
def fail_derivation(*_):
|
||||||
|
raise RuntimeError("derivation still unavailable")
|
||||||
|
|
||||||
|
monkeypatch.setattr(collector, "derive_hours_from_interval", fail_derivation)
|
||||||
|
failed = collector.run_collection(
|
||||||
|
config=config,
|
||||||
|
clock=FakeClock(now),
|
||||||
|
acquire=_acquirer(sysfs_path, 1010),
|
||||||
|
)
|
||||||
|
assert failed["ok"] is False
|
||||||
|
assert "capacity full (1 observations)" in failed["error"]
|
||||||
|
|
||||||
|
called = False
|
||||||
|
|
||||||
|
def must_not_acquire():
|
||||||
|
nonlocal called
|
||||||
|
called = True
|
||||||
|
return _smartctl(1020), sysfs_path
|
||||||
|
|
||||||
|
retry = collector.run_collection(
|
||||||
|
config=config,
|
||||||
|
clock=FakeClock(now + timedelta(minutes=5)),
|
||||||
|
acquire=must_not_acquire,
|
||||||
|
)
|
||||||
|
assert retry["ok"] is False
|
||||||
|
assert "capacity full (1 observations)" in retry["error"]
|
||||||
|
assert called is False
|
||||||
|
|
||||||
|
|
||||||
|
def test_v3_readers_ignore_pending_table_until_store_migrates(tmp_path):
|
||||||
|
store_path = tmp_path / "observations.db"
|
||||||
|
conn = init_store(store_path)
|
||||||
|
conn.execute("DROP TABLE pending_publications")
|
||||||
|
conn.execute("PRAGMA user_version=3")
|
||||||
|
conn.commit()
|
||||||
|
conn.close()
|
||||||
|
|
||||||
|
with read_status(store_path, datetime.now(timezone.utc), query_services=False) as (reader, comp):
|
||||||
|
assert reader is not None
|
||||||
|
assert comp.store_fault is None
|
||||||
|
assert comp.pending_publication_count == 0
|
||||||
|
|
||||||
|
migrated = init_store(store_path)
|
||||||
|
try:
|
||||||
|
assert migrated.execute("PRAGMA user_version").fetchone()[0] == 4
|
||||||
|
assert migrated.execute("SELECT COUNT(*) FROM pending_publications").fetchone()[0] == 0
|
||||||
|
finally:
|
||||||
|
migrated.close()
|
||||||
Reference in New Issue
Block a user