Compare commits
5
Commits
v0.5.0
...
8a5ef05188
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
8a5ef05188 | ||
|
|
847bde1be0 | ||
|
|
0d45f263d2 | ||
|
|
7c21b044ba | ||
|
|
437ea6be73 |
+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)")
|
||||||
|
|||||||
+188
-94
@@ -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
|
||||||
- Commit one well-formed sample
|
- 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.
|
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,12 +200,12 @@ 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.
|
||||||
|
|
||||||
Identity normalization happens exactly once here.
|
Identity normalization happens exactly once here.
|
||||||
Returns segment info for the caller.
|
Returns segment info for the caller. Caller owns the transaction.
|
||||||
"""
|
"""
|
||||||
from .segment import find_current_segment, should_open_new_segment, open_segment
|
from .segment import find_current_segment, should_open_new_segment, open_segment
|
||||||
|
|
||||||
@@ -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
|
||||||
@@ -262,8 +266,6 @@ def write_sample(
|
|||||||
),
|
),
|
||||||
)
|
)
|
||||||
|
|
||||||
conn.commit()
|
|
||||||
|
|
||||||
return {
|
return {
|
||||||
"segment_opened": segment_opened,
|
"segment_opened": segment_opened,
|
||||||
"segment_reason": reason,
|
"segment_reason": reason,
|
||||||
@@ -273,108 +275,200 @@ def write_sample(
|
|||||||
}
|
}
|
||||||
|
|
||||||
|
|
||||||
def run_collection(
|
def _observation_time(sample: Dict[str, Any]) -> datetime:
|
||||||
smartctl_data: Dict[str, Any],
|
"""Read an observation's original timestamp for ordered recovery."""
|
||||||
sysfs_path: Path,
|
observed_at = datetime.fromisoformat(sample["ts"])
|
||||||
config: Dict[str, Any],
|
if observed_at.tzinfo is None:
|
||||||
clock,
|
return observed_at.replace(tzinfo=timezone.utc)
|
||||||
) -> Dict[str, Any]:
|
return observed_at.astimezone(timezone.utc)
|
||||||
"""Run one collection run.
|
|
||||||
|
|
||||||
This is the main entry point for the collector.
|
|
||||||
Returns the run outcome.
|
|
||||||
"""
|
|
||||||
try:
|
|
||||||
# Acquire counters and thermal evidence
|
|
||||||
counters = acquire_from_smartctl(smartctl_data)
|
|
||||||
|
|
||||||
# Acquire controller identity
|
def _publish_observation(
|
||||||
identity = acquire_from_sysfs(sysfs_path)
|
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)
|
||||||
|
|
||||||
# Build sample with injected clock
|
cursor = conn.execute("SELECT id FROM samples ORDER BY id DESC LIMIT 1")
|
||||||
sample = {
|
current_id = cursor.fetchone()[0]
|
||||||
"ts": clock.utcnow().isoformat(),
|
prev = find_previous_sample(conn, seg_info.get("segment_id"), current_id)
|
||||||
"device": config["device"],
|
if prev is not None:
|
||||||
**counters,
|
current = {
|
||||||
**identity,
|
"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(
|
||||||
|
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")
|
||||||
|
|
||||||
# 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)
|
||||||
|
|
||||||
|
published_count = _recover_pending_for_collection(conn)
|
||||||
|
|
||||||
try:
|
# Hold the writer reservation across the capacity check and acquisition.
|
||||||
# Ensure monitoring period is open (issue #73 AC2)
|
# A concurrent collector will recheck pending work before it acquires.
|
||||||
ensure_period_open(conn, clock.utcnow())
|
while True:
|
||||||
|
conn.execute("BEGIN IMMEDIATE")
|
||||||
|
waiting = _pending_count(conn)
|
||||||
|
if waiting:
|
||||||
|
conn.rollback()
|
||||||
|
published_count += _recover_pending_for_collection(conn)
|
||||||
|
continue
|
||||||
|
|
||||||
# Validate invariants
|
from .tz_util import detect_system_tz
|
||||||
validate_sample_invariants(sample, conn)
|
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")
|
||||||
|
|
||||||
# Write sample and get segment info
|
counters = acquire_from_smartctl(smartctl_data)
|
||||||
seg_info = write_sample(sample, identity, conn, clock)
|
identity = acquire_from_sysfs(sysfs_path)
|
||||||
|
sample = {
|
||||||
# Derive hour observations from interval with previous sample
|
"ts": clock.utcnow().isoformat(),
|
||||||
try:
|
"device": config["device"],
|
||||||
# Find the sample we just wrote
|
**counters,
|
||||||
cursor = conn.execute("SELECT id FROM samples ORDER BY id DESC LIMIT 1")
|
**identity,
|
||||||
current_id = cursor.fetchone()[0]
|
|
||||||
|
|
||||||
prev = find_previous_sample(conn, seg_info.get("segment_id"), current_id)
|
|
||||||
if prev is not None:
|
|
||||||
# Build current sample dict for derivation
|
|
||||||
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)
|
|
||||||
|
|
||||||
# Rebuild day aggregates from hour observations
|
|
||||||
from .day_aggregate import derive_all_days, persist_day_aggregate
|
|
||||||
for agg in derive_all_days(conn):
|
|
||||||
persist_day_aggregate(conn, agg)
|
|
||||||
|
|
||||||
# Derive local-day summary using system timezone (issue #90)
|
|
||||||
try:
|
|
||||||
from .tz_util import detect_system_tz
|
|
||||||
from .local_day import derive_local_day_summary, persist_local_day
|
|
||||||
tz_name = detect_system_tz()
|
|
||||||
clock_now = clock.utcnow()
|
|
||||||
local_summary = derive_local_day_summary(conn, tz_name, clock_now)
|
|
||||||
if local_summary is not None:
|
|
||||||
persist_local_day(conn, local_summary)
|
|
||||||
except Exception:
|
|
||||||
# Local-day derivation failure must not prevent publication
|
|
||||||
pass
|
|
||||||
|
|
||||||
conn.commit()
|
|
||||||
except Exception:
|
|
||||||
# Derivation failure must not prevent sample persistence (issue #73 AC6)
|
|
||||||
pass
|
|
||||||
|
|
||||||
return {
|
|
||||||
"ok": True,
|
|
||||||
"sample_count": 1,
|
|
||||||
"store_path": str(store_path),
|
|
||||||
}
|
}
|
||||||
finally:
|
validate_sample_invariants(sample, conn)
|
||||||
conn.close()
|
_stage_observation(conn, sample, identity, tz_name)
|
||||||
|
conn.commit()
|
||||||
|
break
|
||||||
|
|
||||||
except (AcquisitionError, InvariantViolationError) as e:
|
published_count += _recover_pending_for_collection(conn)
|
||||||
|
return {
|
||||||
|
"ok": True,
|
||||||
|
"sample_count": published_count,
|
||||||
|
"store_path": str(store_path),
|
||||||
|
}
|
||||||
|
|
||||||
|
except Exception as exc:
|
||||||
|
if conn is not None:
|
||||||
|
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:
|
||||||
|
if conn is not None:
|
||||||
|
conn.close()
|
||||||
|
|||||||
@@ -96,6 +96,7 @@ def derive_hours_from_interval(
|
|||||||
"""Derive hour observations from a sample pair interval.
|
"""Derive hour observations from a sample pair interval.
|
||||||
|
|
||||||
Returns list of hour observation dicts that were written/updated.
|
Returns list of hour observation dicts that were written/updated.
|
||||||
|
Caller owns the transaction.
|
||||||
"""
|
"""
|
||||||
prev_ts = _parse_ts(prev_sample["ts"])
|
prev_ts = _parse_ts(prev_sample["ts"])
|
||||||
next_ts = _parse_ts(next_sample["ts"])
|
next_ts = _parse_ts(next_sample["ts"])
|
||||||
@@ -207,9 +208,6 @@ def _upsert_hour_observation(
|
|||||||
"WHERE id = ?",
|
"WHERE id = ?",
|
||||||
(new_bw, new_br, new_samples, existing[0]),
|
(new_bw, new_br, new_samples, existing[0]),
|
||||||
)
|
)
|
||||||
conn.commit()
|
|
||||||
|
|
||||||
|
|
||||||
def _add_unattributed_bytes(
|
def _add_unattributed_bytes(
|
||||||
conn: sqlite3.Connection,
|
conn: sqlite3.Connection,
|
||||||
prev_ts: datetime,
|
prev_ts: datetime,
|
||||||
@@ -241,4 +239,3 @@ def _add_unattributed_bytes(
|
|||||||
"unattributed_bytes_read = unattributed_bytes_read + ? WHERE day = ?",
|
"unattributed_bytes_read = unattributed_bytes_read + ? WHERE day = ?",
|
||||||
(bw_delta, br_delta, day),
|
(bw_delta, br_delta, day),
|
||||||
)
|
)
|
||||||
conn.commit()
|
|
||||||
|
|||||||
@@ -69,6 +69,7 @@ def cmd_enable(args: argparse.Namespace) -> None:
|
|||||||
open_period = get_open_period(conn)
|
open_period = get_open_period(conn)
|
||||||
if open_period is None:
|
if open_period is None:
|
||||||
ensure_period_open(conn, now)
|
ensure_period_open(conn, now)
|
||||||
|
conn.commit()
|
||||||
print("Monitoring period opened at", now.isoformat())
|
print("Monitoring period opened at", now.isoformat())
|
||||||
else:
|
else:
|
||||||
print("Monitoring period already open (id=%d)" % open_period["id"])
|
print("Monitoring period already open (id=%d)" % open_period["id"])
|
||||||
|
|||||||
@@ -16,7 +16,7 @@ def ensure_period_open(conn: sqlite3.Connection, run_time: datetime) -> None:
|
|||||||
"""Ensure a monitoring period is open. If none exists, open one at run_time.
|
"""Ensure a monitoring period is open. If none exists, open one at run_time.
|
||||||
|
|
||||||
Spec §9.8: A collection run finding no open monitoring period opens one
|
Spec §9.8: A collection run finding no open monitoring period opens one
|
||||||
at the run moment, never backdated.
|
at the run moment, never backdated. Caller owns the transaction.
|
||||||
"""
|
"""
|
||||||
if get_open_period(conn) is not None:
|
if get_open_period(conn) is not None:
|
||||||
return # Already open — no-op
|
return # Already open — no-op
|
||||||
@@ -26,7 +26,6 @@ def ensure_period_open(conn: sqlite3.Connection, run_time: datetime) -> None:
|
|||||||
"INSERT INTO monitoring_periods (started_at) VALUES (?)",
|
"INSERT INTO monitoring_periods (started_at) VALUES (?)",
|
||||||
(ts,),
|
(ts,),
|
||||||
)
|
)
|
||||||
conn.commit()
|
|
||||||
|
|
||||||
|
|
||||||
def close_period(
|
def close_period(
|
||||||
|
|||||||
@@ -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.
|
||||||
|
|||||||
@@ -25,11 +25,14 @@ def detect_system_tz() -> str:
|
|||||||
|
|
||||||
localtime = Path("/etc/localtime")
|
localtime = Path("/etc/localtime")
|
||||||
if localtime.is_symlink():
|
if localtime.is_symlink():
|
||||||
target = os.readlink(str(localtime))
|
target_path = Path(os.readlink(str(localtime)))
|
||||||
# Strip common prefixes: /usr/share/zoneinfo/, /usr/lib/zoneinfo/
|
if not target_path.is_absolute():
|
||||||
for prefix in ("/usr/share/zoneinfo/", "/usr/lib/zoneinfo/"):
|
target_path = localtime.parent / target_path
|
||||||
if target.startswith(prefix):
|
target = target_path.resolve().as_posix()
|
||||||
return target[len(prefix):]
|
marker = "/zoneinfo/"
|
||||||
|
marker_index = target.find(marker)
|
||||||
|
if marker_index >= 0:
|
||||||
|
return target[marker_index + len(marker):]
|
||||||
return target
|
return target
|
||||||
|
|
||||||
return "UTC"
|
return "UTC"
|
||||||
|
|||||||
@@ -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))
|
||||||
@@ -356,6 +356,165 @@ class TestCollectorDerivation:
|
|||||||
assert count_after == 2
|
assert count_after == 2
|
||||||
|
|
||||||
|
|
||||||
|
class TestCollectionAtomicity:
|
||||||
|
"""Collection exposes sample and derived evidence as one publication."""
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_collection_is_visible_through_cli_and_tui_readers(
|
||||||
|
self, config_fixture, smartctl_fixture, sysfs_fixture_tree, monkeypatch
|
||||||
|
):
|
||||||
|
"""Ordinary readers observe matching published sample and activity evidence."""
|
||||||
|
from fenris.status import get_status, read_status
|
||||||
|
from fenris.tui import FenrisTuiApp
|
||||||
|
|
||||||
|
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)
|
||||||
|
first = {
|
||||||
|
**smartctl_fixture,
|
||||||
|
"nvme_smart_health_information_log": {
|
||||||
|
**smartctl_fixture["nvme_smart_health_information_log"],
|
||||||
|
"data_units_written": 12345678,
|
||||||
|
"data_units_read": 9876543,
|
||||||
|
},
|
||||||
|
}
|
||||||
|
second = {
|
||||||
|
**smartctl_fixture,
|
||||||
|
"nvme_smart_health_information_log": {
|
||||||
|
**smartctl_fixture["nvme_smart_health_information_log"],
|
||||||
|
"data_units_written": 12345698,
|
||||||
|
"data_units_read": 9876553,
|
||||||
|
},
|
||||||
|
}
|
||||||
|
sysfs_path = sysfs_fixture_tree / "sys" / "class" / "nvme" / "nvme0"
|
||||||
|
|
||||||
|
assert run_collection(
|
||||||
|
first, sysfs_path, config_fixture, FakeClock(now - timedelta(minutes=5))
|
||||||
|
)["ok"] is True
|
||||||
|
assert run_collection(
|
||||||
|
second, sysfs_path, config_fixture, FakeClock(now)
|
||||||
|
)["ok"] is True
|
||||||
|
|
||||||
|
store_path = Path(config_fixture["store_path"])
|
||||||
|
with read_status(store_path, now, query_services=False) as (reader, composition):
|
||||||
|
assert reader is not None
|
||||||
|
assert composition.sample_count == 2
|
||||||
|
assert composition.day_count == 1
|
||||||
|
utc_day = reader.execute(
|
||||||
|
"SELECT bytes_written_delta, bytes_read_delta FROM day_aggregates"
|
||||||
|
).fetchone()
|
||||||
|
local_day = reader.execute(
|
||||||
|
"SELECT bytes_written, bytes_read FROM local_days"
|
||||||
|
).fetchone()
|
||||||
|
assert tuple(utc_day) == (10_240_000, 5_120_000)
|
||||||
|
assert tuple(local_day) == (10_240_000, 5_120_000)
|
||||||
|
|
||||||
|
cli_output = get_status(
|
||||||
|
store_path, now, query_services=True, query_journal=False
|
||||||
|
)
|
||||||
|
assert "Monitoring" in cli_output
|
||||||
|
|
||||||
|
app = FenrisTuiApp(store_path=store_path)
|
||||||
|
async with app.run_test(size=(100, 30)) as pilot:
|
||||||
|
await pilot.pause()
|
||||||
|
live_readout = str(app.query_one("#live-readout").render())
|
||||||
|
assert "W 0.010 GB" in live_readout
|
||||||
|
assert "R 0.005 GB" in live_readout
|
||||||
|
|
||||||
|
repeated_result = run_collection(
|
||||||
|
second, sysfs_path, config_fixture, FakeClock(now + timedelta(minutes=5))
|
||||||
|
)
|
||||||
|
assert repeated_result["ok"] is True
|
||||||
|
with read_status(
|
||||||
|
store_path, now + timedelta(minutes=5), query_services=False
|
||||||
|
) as (reader, composition):
|
||||||
|
assert reader is not None
|
||||||
|
assert composition.sample_count == 3
|
||||||
|
utc_day = reader.execute(
|
||||||
|
"SELECT bytes_written_delta, bytes_read_delta FROM day_aggregates"
|
||||||
|
).fetchone()
|
||||||
|
local_day = reader.execute(
|
||||||
|
"SELECT bytes_written, bytes_read FROM local_days"
|
||||||
|
).fetchone()
|
||||||
|
assert tuple(utc_day) == (10_240_000, 5_120_000)
|
||||||
|
assert tuple(local_day) == (10_240_000, 5_120_000)
|
||||||
|
|
||||||
|
def test_failed_local_day_publication_keeps_previous_publication(
|
||||||
|
self, config_fixture, smartctl_fixture, sysfs_fixture_tree, monkeypatch
|
||||||
|
):
|
||||||
|
"""A failed final derivation step leaves all prior reader state intact."""
|
||||||
|
from fenris.status import read_status
|
||||||
|
|
||||||
|
monkeypatch.setenv("TZ", "UTC")
|
||||||
|
now = datetime.now(timezone.utc).replace(second=0, microsecond=0)
|
||||||
|
first_clock = FakeClock(now)
|
||||||
|
first = {
|
||||||
|
**smartctl_fixture,
|
||||||
|
"nvme_smart_health_information_log": {
|
||||||
|
**smartctl_fixture["nvme_smart_health_information_log"],
|
||||||
|
"data_units_written": 12345678,
|
||||||
|
"data_units_read": 9876543,
|
||||||
|
},
|
||||||
|
}
|
||||||
|
first_result = run_collection(
|
||||||
|
first,
|
||||||
|
sysfs_fixture_tree / "sys" / "class" / "nvme" / "nvme0",
|
||||||
|
config_fixture,
|
||||||
|
first_clock,
|
||||||
|
)
|
||||||
|
assert first_result["ok"] is True, first_result
|
||||||
|
|
||||||
|
writer = sqlite3.connect(config_fixture["store_path"])
|
||||||
|
writer.execute(
|
||||||
|
"CREATE TRIGGER fail_local_day_publication "
|
||||||
|
"BEFORE INSERT ON local_days "
|
||||||
|
"BEGIN SELECT RAISE(ABORT, 'injected local-day publication failure'); END"
|
||||||
|
)
|
||||||
|
writer.commit()
|
||||||
|
writer.close()
|
||||||
|
|
||||||
|
next_time = now + timedelta(minutes=5)
|
||||||
|
second = {
|
||||||
|
**smartctl_fixture,
|
||||||
|
"nvme_smart_health_information_log": {
|
||||||
|
**smartctl_fixture["nvme_smart_health_information_log"],
|
||||||
|
"data_units_written": 12345698,
|
||||||
|
"data_units_read": 9876548,
|
||||||
|
},
|
||||||
|
}
|
||||||
|
failed_result = run_collection(
|
||||||
|
second,
|
||||||
|
sysfs_fixture_tree / "sys" / "class" / "nvme" / "nvme0",
|
||||||
|
config_fixture,
|
||||||
|
FakeClock(next_time),
|
||||||
|
)
|
||||||
|
|
||||||
|
assert failed_result["ok"] is False
|
||||||
|
assert "injected local-day publication failure" in failed_result["error"]
|
||||||
|
|
||||||
|
with read_status(
|
||||||
|
Path(config_fixture["store_path"]), next_time, query_services=False
|
||||||
|
) as (reader, composition):
|
||||||
|
assert reader is not None
|
||||||
|
assert composition.sample_count == 1
|
||||||
|
public_counts = reader.execute(
|
||||||
|
"SELECT (SELECT COUNT(*) FROM samples), "
|
||||||
|
"(SELECT COUNT(*) FROM controller_segments), "
|
||||||
|
"(SELECT COUNT(*) FROM monitoring_periods), "
|
||||||
|
"(SELECT COUNT(*) FROM hour_observations), "
|
||||||
|
"(SELECT COUNT(*) FROM day_aggregates), "
|
||||||
|
"(SELECT COUNT(*) FROM local_days)"
|
||||||
|
).fetchone()
|
||||||
|
|
||||||
|
assert tuple(public_counts) == (1, 1, 1, 0, 0, 0)
|
||||||
|
|
||||||
|
|
||||||
# ---------------------------------------------------------------------------
|
# ---------------------------------------------------------------------------
|
||||||
# Display states: awaiting first sample, awaiting another sample
|
# Display states: awaiting first sample, awaiting another sample
|
||||||
# ---------------------------------------------------------------------------
|
# ---------------------------------------------------------------------------
|
||||||
|
|||||||
@@ -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()
|
||||||
|
|||||||
@@ -634,6 +634,14 @@ class TestMonitorIntegration:
|
|||||||
cmd_enable(args)
|
cmd_enable(args)
|
||||||
mock_enable.assert_called_once_with(True)
|
mock_enable.assert_called_once_with(True)
|
||||||
|
|
||||||
|
conn = sqlite3.connect(store_path)
|
||||||
|
try:
|
||||||
|
assert conn.execute(
|
||||||
|
"SELECT COUNT(*) FROM monitoring_periods WHERE ended_at IS NULL"
|
||||||
|
).fetchone()[0] == 1
|
||||||
|
finally:
|
||||||
|
conn.close()
|
||||||
|
|
||||||
def test_cmd_disable_uses_init_system(self, tmp_path):
|
def test_cmd_disable_uses_init_system(self, tmp_path):
|
||||||
"""cmd_disable uses init_system.disable_timer."""
|
"""cmd_disable uses init_system.disable_timer."""
|
||||||
from fenris.monitor import cmd_disable
|
from fenris.monitor import cmd_disable
|
||||||
|
|||||||
@@ -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()
|
||||||
@@ -109,6 +109,7 @@ def _insert_sample(conn, ts, pu=5, device="/dev/nvme0n1"):
|
|||||||
|
|
||||||
def _open_period(conn, start="2026-09-01T00:00:00+00:00"):
|
def _open_period(conn, start="2026-09-01T00:00:00+00:00"):
|
||||||
ensure_period_open(conn, datetime.fromisoformat(start))
|
ensure_period_open(conn, datetime.fromisoformat(start))
|
||||||
|
conn.commit()
|
||||||
|
|
||||||
|
|
||||||
def _insert_local_day(conn, local_date, tz_name="UTC", tz_offset="+00:00",
|
def _insert_local_day(conn, local_date, tz_name="UTC", tz_offset="+00:00",
|
||||||
|
|||||||
@@ -0,0 +1,26 @@
|
|||||||
|
"""Timezone detection tests."""
|
||||||
|
import sys
|
||||||
|
from pathlib import Path
|
||||||
|
|
||||||
|
sys.path.insert(0, str(Path(__file__).parent.parent / "src"))
|
||||||
|
|
||||||
|
from fenris import tz_util
|
||||||
|
|
||||||
|
|
||||||
|
def test_relative_localtime_symlink_returns_zoneinfo_key(tmp_path, monkeypatch):
|
||||||
|
zoneinfo = tmp_path / "usr" / "share" / "zoneinfo" / "Asia" / "Kolkata"
|
||||||
|
zoneinfo.parent.mkdir(parents=True)
|
||||||
|
zoneinfo.write_bytes(b"zoneinfo")
|
||||||
|
localtime = tmp_path / "etc" / "localtime"
|
||||||
|
localtime.parent.mkdir()
|
||||||
|
localtime.symlink_to("../usr/share/zoneinfo/Asia/Kolkata")
|
||||||
|
|
||||||
|
real_path = tz_util.Path
|
||||||
|
monkeypatch.setattr(
|
||||||
|
tz_util,
|
||||||
|
"Path",
|
||||||
|
lambda path: localtime if path == "/etc/localtime" else real_path(path),
|
||||||
|
)
|
||||||
|
monkeypatch.delenv("TZ", raising=False)
|
||||||
|
|
||||||
|
assert tz_util.detect_system_tz() == "Asia/Kolkata"
|
||||||
Reference in New Issue
Block a user