Compare commits
3
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
8a5ef05188 | ||
|
|
847bde1be0 | ||
|
|
0d45f263d2 |
+9
-10
@@ -27,7 +27,6 @@ if VENV_DIR.exists():
|
||||
if site_packages:
|
||||
sys.path.insert(0, str(site_packages))
|
||||
|
||||
from fenris.store import init_store, get_store_path
|
||||
from fenris.collector import run_collection
|
||||
|
||||
|
||||
@@ -103,14 +102,6 @@ def main() -> None:
|
||||
config = load_config()
|
||||
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
|
||||
class SimpleClock:
|
||||
def utcnow(self):
|
||||
@@ -118,8 +109,16 @@ def main() -> None:
|
||||
|
||||
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
|
||||
result = run_collection(smartctl_data, sysfs_path, config, clock)
|
||||
result = run_collection(config=config, clock=clock, acquire=acquire)
|
||||
|
||||
if result["ok"]:
|
||||
print(f"Collection successful: {result['sample_count']} sample(s)")
|
||||
|
||||
+167
-55
@@ -5,19 +5,21 @@ This module implements the thinnest complete write path:
|
||||
- Acquire controller identity from sysfs
|
||||
- Normalize identity exactly once at write time
|
||||
- Validate every row against store invariants
|
||||
- Stage valid observations privately before derivation
|
||||
- Publish the sample and derived evidence in one collection-owned transaction
|
||||
|
||||
No code path outside the collector interrogates the device.
|
||||
"""
|
||||
import json
|
||||
import sqlite3
|
||||
from collections.abc import Callable
|
||||
from datetime import datetime, timezone
|
||||
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 .derive import find_previous_sample, derive_hours_from_interval
|
||||
from .store import init_store, get_store_path
|
||||
|
||||
|
||||
class AcquisitionError(Exception):
|
||||
@@ -30,6 +32,9 @@ class InvariantViolationError(Exception):
|
||||
pass
|
||||
|
||||
|
||||
PENDING_PUBLICATION_LIMIT = 6720
|
||||
|
||||
|
||||
def acquire_from_smartctl(smartctl_data: Dict[str, Any]) -> Dict[str, Any]:
|
||||
"""Acquire counters and thermal evidence from smartctl -a -j data.
|
||||
|
||||
@@ -195,7 +200,7 @@ def write_sample(
|
||||
sample: Dict[str, Any],
|
||||
identity: Dict[str, Any],
|
||||
conn: sqlite3.Connection,
|
||||
clock,
|
||||
observed_at: datetime,
|
||||
) -> Dict[str, Any]:
|
||||
"""Write one sample to the observation store.
|
||||
|
||||
@@ -219,8 +224,7 @@ def write_sample(
|
||||
# Open new segment if needed
|
||||
segment_opened = False
|
||||
if should_open:
|
||||
now = clock.utcnow()
|
||||
open_segment(conn, now, identity, identity_key, identity_degraded)
|
||||
open_segment(conn, observed_at, identity, identity_key, identity_degraded)
|
||||
segment_opened = True
|
||||
|
||||
# Get current segment_id for provenance
|
||||
@@ -271,48 +275,25 @@ def write_sample(
|
||||
}
|
||||
|
||||
|
||||
def run_collection(
|
||||
smartctl_data: Dict[str, Any],
|
||||
sysfs_path: Path,
|
||||
config: Dict[str, Any],
|
||||
clock,
|
||||
) -> Dict[str, Any]:
|
||||
"""Run one collection run.
|
||||
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)
|
||||
|
||||
This is the main entry point for the collector.
|
||||
Returns the run outcome.
|
||||
"""
|
||||
conn = None
|
||||
try:
|
||||
# Acquire counters and thermal evidence
|
||||
counters = acquire_from_smartctl(smartctl_data)
|
||||
|
||||
# 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)
|
||||
conn = init_store(store_path)
|
||||
|
||||
# Run legacy import if needed (idempotent)
|
||||
from .legacy import import_legacy_history
|
||||
history_path = Path(config.get("data_dir", ".")) / "history.jsonl"
|
||||
if history_path.exists():
|
||||
import_legacy_history(conn, history_path, clock=clock)
|
||||
|
||||
# Collection owns one transaction for the sample and its evidence.
|
||||
ensure_period_open(conn, clock.utcnow())
|
||||
|
||||
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)
|
||||
seg_info = write_sample(sample, identity, conn, clock)
|
||||
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]
|
||||
@@ -331,31 +312,162 @@ def run_collection(
|
||||
derive_hours_from_interval(conn, prev, current)
|
||||
|
||||
from .day_aggregate import derive_all_days, persist_day_aggregate
|
||||
for agg in derive_all_days(conn):
|
||||
persist_day_aggregate(conn, agg)
|
||||
for aggregate in derive_all_days(conn):
|
||||
persist_day_aggregate(conn, aggregate)
|
||||
|
||||
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())
|
||||
local_summary = derive_local_day_summary(conn, tz_name, observed_at)
|
||||
if local_summary is not None:
|
||||
persist_local_day(conn, local_summary)
|
||||
|
||||
conn.commit()
|
||||
|
||||
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(
|
||||
smartctl_data: Optional[Dict[str, Any]] = None,
|
||||
sysfs_path: Optional[Path] = None,
|
||||
config: Optional[Dict[str, Any]] = None,
|
||||
clock=None,
|
||||
*,
|
||||
acquire: Optional[Callable[[], tuple[Dict[str, Any], Path]]] = None,
|
||||
) -> Dict[str, Any]:
|
||||
"""Recover old work, acquire one observation, then publish it atomically.
|
||||
|
||||
Production callers pass ``acquire`` so recovery and capacity checks run
|
||||
before device interrogation. Direct sample arguments remain useful for
|
||||
deterministic collector tests.
|
||||
"""
|
||||
conn = None
|
||||
try:
|
||||
if config is None or clock is None:
|
||||
raise ValueError("config and clock are required")
|
||||
|
||||
store_path = get_store_path(config)
|
||||
conn = init_store(store_path)
|
||||
|
||||
from .legacy import import_legacy_history
|
||||
history_path = Path(config.get("data_dir", ".")) / "history.jsonl"
|
||||
if history_path.exists():
|
||||
import_legacy_history(conn, history_path, clock=clock)
|
||||
|
||||
published_count = _recover_pending_for_collection(conn)
|
||||
|
||||
# Hold the writer reservation across the capacity check and acquisition.
|
||||
# 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
|
||||
|
||||
from .tz_util import detect_system_tz
|
||||
tz_name = detect_system_tz()
|
||||
if acquire is not None:
|
||||
smartctl_data, sysfs_path = acquire()
|
||||
if smartctl_data is None or sysfs_path is None:
|
||||
raise AcquisitionError("No acquired SMART data or sysfs identity")
|
||||
|
||||
counters = acquire_from_smartctl(smartctl_data)
|
||||
identity = acquire_from_sysfs(sysfs_path)
|
||||
sample = {
|
||||
"ts": clock.utcnow().isoformat(),
|
||||
"device": config["device"],
|
||||
**counters,
|
||||
**identity,
|
||||
}
|
||||
validate_sample_invariants(sample, conn)
|
||||
_stage_observation(conn, sample, identity, tz_name)
|
||||
conn.commit()
|
||||
break
|
||||
|
||||
published_count += _recover_pending_for_collection(conn)
|
||||
return {
|
||||
"ok": True,
|
||||
"sample_count": 1,
|
||||
"sample_count": published_count,
|
||||
"store_path": str(store_path),
|
||||
}
|
||||
|
||||
except Exception as e:
|
||||
except Exception as exc:
|
||||
if conn is not None:
|
||||
conn.rollback()
|
||||
return {
|
||||
"ok": False,
|
||||
"error": str(e),
|
||||
"error_type": type(e).__name__,
|
||||
"error": str(exc),
|
||||
"error_type": type(exc).__name__,
|
||||
}
|
||||
finally:
|
||||
if conn is not None:
|
||||
|
||||
@@ -115,6 +115,7 @@ class StatusComposition:
|
||||
# Sample counts for waiting explanations
|
||||
sample_count: int = 0
|
||||
day_count: int = 0
|
||||
pending_publication_count: int = 0
|
||||
|
||||
# Whether the status dot should blink (only Monitoring)
|
||||
should_blink: bool = False
|
||||
@@ -286,6 +287,7 @@ def compose_status(
|
||||
freshness_age_s = None
|
||||
sample_count = 0
|
||||
day_count = 0
|
||||
pending_publication_count = 0
|
||||
deliberately_paused = False
|
||||
|
||||
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]
|
||||
cursor = conn.execute("SELECT COUNT(*) FROM day_aggregates")
|
||||
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:
|
||||
store_fault = str(exc)
|
||||
|
||||
@@ -328,6 +334,7 @@ def compose_status(
|
||||
freshness = "unknown"
|
||||
freshness_age_s = None
|
||||
sample_count = day_count = 0
|
||||
pending_publication_count = 0
|
||||
deliberately_paused = False
|
||||
|
||||
# --- External stop detection ---
|
||||
@@ -417,6 +424,7 @@ def compose_status(
|
||||
overlay_base=overlay_base,
|
||||
sample_count=sample_count,
|
||||
day_count=day_count,
|
||||
pending_publication_count=pending_publication_count,
|
||||
continuity=continuity,
|
||||
paused_lines=paused_lines,
|
||||
)
|
||||
@@ -459,6 +467,8 @@ def render_status_cli(comp: StatusComposition) -> str:
|
||||
facts.append("timer: %s" % (
|
||||
"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:
|
||||
lines.append(" · ".join(facts))
|
||||
@@ -524,6 +534,8 @@ def render_status_tui(comp: StatusComposition) -> str:
|
||||
facts.append("Timer: %s" % (
|
||||
"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:
|
||||
lines.append(" · ".join(facts))
|
||||
@@ -539,3 +551,8 @@ def render_status_tui(comp: StatusComposition) -> str:
|
||||
lines.append(pl[:1].upper() + pl[1:])
|
||||
|
||||
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:
|
||||
- Store initialization with WAL mode
|
||||
- Schema versioning with PRAGMA user_version
|
||||
- The six entities: samples, hour_observations, day_aggregates,
|
||||
monitoring_periods, controller_segments, endurance_baseline
|
||||
- Observation records, derived activity, publication state, and metadata
|
||||
"""
|
||||
import sqlite3
|
||||
from pathlib import Path
|
||||
@@ -12,7 +11,7 @@ from typing import Optional
|
||||
|
||||
|
||||
# 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
|
||||
@@ -33,7 +32,7 @@ def get_store_path(config: dict) -> Path:
|
||||
def init_store(store_path: Path) -> sqlite3.Connection:
|
||||
"""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.
|
||||
"""
|
||||
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):
|
||||
"""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
|
||||
|
||||
# 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:
|
||||
"""Apply forward-only migrations to bring the store to SCHEMA_VERSION.
|
||||
|
||||
@@ -154,7 +154,7 @@ class TestSchemaMigration:
|
||||
|
||||
# Migrate
|
||||
steps = migrate_to_latest(db)
|
||||
assert steps == 2 # v1→v2→v3
|
||||
assert steps == 3 # v1→v2→v3→v4
|
||||
|
||||
# Verify data preserved
|
||||
conn = sqlite3.connect(str(db))
|
||||
|
||||
@@ -232,7 +232,7 @@ def test_store_initialization(config_fixture: Dict[str, Any]):
|
||||
# Initialize store
|
||||
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'")
|
||||
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("store_metadata")
|
||||
expected_tables.add("local_days")
|
||||
expected_tables.add("pending_publications")
|
||||
assert expected_tables == tables
|
||||
|
||||
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