diff --git a/src/fenris/collector.py b/src/fenris/collector.py index 8f5e16d..4c90781 100644 --- a/src/fenris/collector.py +++ b/src/fenris/collector.py @@ -325,7 +325,7 @@ def run_collection( # Find the sample we just wrote 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: # Build current sample dict for derivation @@ -340,6 +340,12 @@ def run_collection( "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) + conn.commit() except Exception: # Derivation failure must not prevent sample persistence (issue #73 AC6) pass diff --git a/src/fenris/day_aggregate.py b/src/fenris/day_aggregate.py index 88f9a80..4c05e6d 100644 --- a/src/fenris/day_aggregate.py +++ b/src/fenris/day_aggregate.py @@ -163,3 +163,41 @@ def derive_all_days(conn: sqlite3.Connection) -> list[DayAggregate]: if agg is not None: results.append(agg) return results + + +def persist_day_aggregate(conn: sqlite3.Connection, agg: DayAggregate) -> None: + """Upsert a derived day aggregate into the day_aggregates table. + + Merges attributed bytes from hour observations with any existing + unattributed cross-hour evidence already stored for this day. + Caller must manage transactions and commits. + """ + existing = conn.execute( + "SELECT id FROM day_aggregates WHERE day = ?", + (agg.day,), + ).fetchone() + + if existing is None: + conn.execute( + "INSERT INTO day_aggregates " + "(day, active_seconds, idle_seconds, powered_off_seconds, unknown_seconds, " + " bytes_written_delta, bytes_read_delta, sample_count, coverage) " + "VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)", + (agg.day, agg.seconds_active, agg.seconds_idle, + agg.seconds_powered_off, agg.seconds_unknown, + agg.bytes_written_delta, agg.bytes_read_delta, + agg.sample_count, agg.coverage), + ) + else: + conn.execute( + "UPDATE day_aggregates " + "SET active_seconds = ?, idle_seconds = ?, powered_off_seconds = ?, " + " unknown_seconds = ?, bytes_written_delta = ?, bytes_read_delta = ?, " + " sample_count = ?, coverage = ? " + "WHERE id = ?", + (agg.seconds_active, agg.seconds_idle, + agg.seconds_powered_off, agg.seconds_unknown, + agg.bytes_written_delta, agg.bytes_read_delta, + agg.sample_count, agg.coverage, + existing[0]), + ) diff --git a/src/fenris/derive.py b/src/fenris/derive.py index 301553b..2470d5f 100644 --- a/src/fenris/derive.py +++ b/src/fenris/derive.py @@ -175,7 +175,8 @@ def _upsert_hour_observation( """Insert or update an hour observation.""" # Check if hour exists existing = conn.execute( - "SELECT id, bytes_written_delta, sample_count FROM hour_observations WHERE hour = ?", + "SELECT id, bytes_written_delta, bytes_read_delta, sample_count " + "FROM hour_observations WHERE hour = ?", (hour_key,), ).fetchone() @@ -198,10 +199,13 @@ def _upsert_hour_observation( else: # Merge: accumulate bytes and sample count new_bw = existing[1] + bw_delta - new_samples = existing[2] + sample_count + new_br = existing[2] + br_delta + new_samples = existing[3] + sample_count conn.execute( - "UPDATE hour_observations SET bytes_written_delta = ?, sample_count = ? WHERE id = ?", - (new_bw, new_samples, existing[0]), + "UPDATE hour_observations " + "SET bytes_written_delta = ?, bytes_read_delta = ?, sample_count = ? " + "WHERE id = ?", + (new_bw, new_br, new_samples, existing[0]), ) conn.commit() diff --git a/src/fenris/tui.py b/src/fenris/tui.py index fa7dea9..fad3ebc 100644 --- a/src/fenris/tui.py +++ b/src/fenris/tui.py @@ -285,15 +285,16 @@ class DailyBarGraph(Widget): return # Summarise visible days as text - total_bytes = sum(d.get("total_bytes", 0) for d in self._day_data) - days_with_data = sum(1 for d in self._day_data if d.get("total_bytes", 0) > 0) + total_written = sum(d.get("total_written", d.get("total_bytes", 0)) for d in self._day_data) + total_read = sum(d.get("total_read", 0) for d in self._day_data) + days_with_data = sum(1 for d in self._day_data if d.get("total_bytes", 0) > 0 or d.get("total_read", 0) > 0) n = len(self._day_data) first = self._day_data[0].get("local_label", "") last = self._day_data[-1].get("local_label", "") self.query_one("#bar-render").update( "[dim]Graph needs ≥80×24[/dim]\n" - " %d days · %d with writes · %.3f GB total\n" - " %s → %s" % (n, days_with_data, total_bytes / 1e9, first, last) + " %d days · %d with activity · W %.3f GB R %.3f GB\n" + " %s → %s" % (n, days_with_data, total_written / 1e9, total_read / 1e9, first, last) ) self.query_one("#bar-legend").update("") @@ -301,30 +302,41 @@ class DailyBarGraph(Widget): if 0 <= self.selected_index < len(self._day_data): day = self._day_data[self.selected_index] self.query_one("#bar-readout").update( - "[bold]%s[/bold] · %.3f GB" - % (day.get("local_label", ""), day.get("total_bytes", 0) / 1e9) + "[bold]%s[/bold] \u00b7 W %.3f GB \u00b7 R %.3f GB" + % ( + day.get("local_label", ""), + day.get("total_written", day.get("total_bytes", 0)) / 1e9, + day.get("total_read", 0) / 1e9, + ) ) else: self.query_one("#bar-readout").update("[dim]No selectable day[/dim]") def _show_constrained_hourly_summary(self) -> None: """Textual fallback for hourly view when terminal is too small.""" - total_bytes = sum(h.get("bytes_written", 0) for h in self._hour_data) + total_written = sum(h.get("bytes_written", 0) for h in self._hour_data) + total_read = sum(h.get("bytes_read", 0) for h in self._hour_data) hours_with_data = sum( - 1 for h in self._hour_data if h.get("bytes_written", 0) > 0 + 1 for h in self._hour_data + if h.get("bytes_written", 0) > 0 or h.get("bytes_read", 0) > 0 ) n = len(self._hour_data) self.query_one("#bar-render").update( "[dim]Graph needs ≥80×24[/dim]\n" - " %d hours · %d with writes · %.3f GB total" % (n, hours_with_data, total_bytes / 1e9) + " %d hours · %d with activity · W %.3f GB R %.3f GB" + % (n, hours_with_data, total_written / 1e9, total_read / 1e9) ) self.query_one("#bar-legend").update("") if 0 <= self._hourly_selected < len(self._hour_data): h = self._hour_data[self._hourly_selected] self.query_one("#bar-readout").update( - "[bold]%s[/bold] \u00b7 %.3f GB" - % (h.get("local_label", ""), h.get("bytes_written", 0) / 1e9) + "[bold]%s[/bold] \u00b7 W %.3f GB \u00b7 R %.3f GB" + % ( + h.get("local_label", ""), + h.get("bytes_written", 0) / 1e9, + h.get("bytes_read", 0) / 1e9, + ) ) else: self.query_one("#bar-readout").update("[dim]\u2190 \u2192 Select hour[/dim]") @@ -461,31 +473,42 @@ class DailyBarGraph(Widget): return day = self._day_data[self.selected_index] - total = day.get("total_bytes", 0) - allocated = day.get("allocated_bytes", 0) - unallocated = day.get("unallocated_bytes", 0) + total_written = day.get("total_written", day.get("total_bytes", 0)) + total_read = day.get("total_read", 0) + allocated_w = day.get("allocated_bytes", 0) + unallocated_w = day.get("unallocated_bytes", 0) + allocated_r = day.get("allocated_read", 0) + unallocated_r = day.get("unallocated_read", 0) coverage = day.get("coverage", 0) hours = day.get("evidenced_hours", 0) - total_gb = total / 1e9 state = " · partial" if day.get("is_partial") else ( " · gap" if day.get("is_gap") else "" ) parts = [ - "[bold]%s UTC[/bold] \u00b7 %.3f GB total \u00b7 %d hours \u00b7 %.0f%% coverage%s" + "[bold]%s UTC[/bold] \u00b7 W %.3f GB \u00b7 R %.3f GB \u00b7 %d hours \u00b7 %.0f%% coverage%s" % ( day.get("local_label", day.get("day", "")), - total_gb, + total_written / 1e9, + total_read / 1e9, hours, coverage * 100, state, ), ] - if unallocated > 0: - parts.append( - " Allocated %.3f GB \u00b7 unallocated %.3f GB" - % (allocated / 1e9, unallocated / 1e9) + alloc_parts = [] + if unallocated_w > 0: + alloc_parts.append( + "W alloc %.3f GB \u00b7 unalloc %.3f GB" + % (allocated_w / 1e9, unallocated_w / 1e9) ) + if unallocated_r > 0: + alloc_parts.append( + "R alloc %.3f GB \u00b7 unalloc %.3f GB" + % (allocated_r / 1e9, unallocated_r / 1e9) + ) + if alloc_parts: + parts.append(" " + " \u00b7 ".join(alloc_parts)) self.query_one("#bar-readout").update("\n".join(parts)) # -- Hourly drill-down rendering -- @@ -580,10 +603,11 @@ class DailyBarGraph(Widget): self._drill_unallocated_bytes / 1e9, ) self.query_one("#bar-readout").update( - "[bold]%s:00 UTC[/bold] \u00b7 %.3f GB \u00b7 %d%% coverage%s%s" + "[bold]%s:00 UTC[/bold] \u00b7 W %.3f GB \u00b7 R %.3f GB \u00b7 %d%% coverage%s%s" % ( h.get("local_label", h.get("hour", "")), h.get("bytes_written", 0) / 1e9, + h.get("bytes_read", 0) / 1e9, h.get("coverage", 0) * 100, state, note, @@ -708,7 +732,8 @@ def _query_daily_graph_data( cursor = conn.execute( "SELECT day, bytes_written_delta, unattributed_bytes_written, " "coverage, sample_count, active_seconds, idle_seconds, " - "powered_off_seconds, unknown_seconds " + "powered_off_seconds, unknown_seconds, " + "bytes_read_delta, unattributed_bytes_read " "FROM day_aggregates ORDER BY day" ) rows = cursor.fetchall() @@ -717,18 +742,21 @@ def _query_daily_graph_data( for row in rows: day = row[0] bw_delta = row[1] or 0 - unattributed = row[2] or 0 + unattributed_w = row[2] or 0 coverage = row[3] or 0.0 sample_count = row[4] or 0 active = row[5] or 0 idle = row[6] or 0 powered_off = row[7] or 0 unknown = row[8] or 0 + br_delta = row[9] or 0 + unattributed_r = row[10] or 0 - total_bytes = bw_delta + unattributed + total_written = bw_delta + unattributed_w + total_read = br_delta + unattributed_r evidenced_hours = (active + idle + powered_off) // 3600 - is_zero = total_bytes == 0 + is_zero = total_written == 0 and total_read == 0 is_gap = ( sample_count == 0 and (active + idle + powered_off) == 0 @@ -739,9 +767,13 @@ def _query_daily_graph_data( by_day[day] = { "day": day, "local_label": day, - "total_bytes": total_bytes, + "total_bytes": total_written, + "total_written": total_written, + "total_read": total_read, "allocated_bytes": bw_delta, - "unallocated_bytes": unattributed, + "unallocated_bytes": unattributed_w, + "allocated_read": br_delta, + "unallocated_read": unattributed_r, "coverage": coverage, "evidenced_hours": evidenced_hours, "sample_count": sample_count, @@ -761,8 +793,12 @@ def _query_daily_graph_data( "day": day, "local_label": day, "total_bytes": 0, + "total_written": 0, + "total_read": 0, "allocated_bytes": 0, "unallocated_bytes": 0, + "allocated_read": 0, + "unallocated_read": 0, "coverage": 0.0, "evidenced_hours": 0, "sample_count": 0, @@ -785,10 +821,10 @@ def _query_hourly_graph_data( ) -> List[Dict[str, Any]]: """Query hour observations for a specific day. - Returns one dict per hour with bytes written, coverage, and flags. + Returns one dict per hour with bytes written/read, coverage, and flags. """ cursor = conn.execute( - "SELECT hour, bytes_written_delta, coverage, sample_count, " + "SELECT hour, bytes_written_delta, bytes_read_delta, coverage, sample_count, " "active_seconds, idle_seconds, powered_off_seconds, unknown_seconds " "FROM hour_observations " "WHERE hour LIKE ? ORDER BY hour", @@ -800,14 +836,15 @@ def _query_hourly_graph_data( for row in rows: hour = row[0] bw = row[1] or 0 - coverage = row[2] or 0.0 - sample_count = row[3] or 0 - active = row[4] or 0 - idle = row[5] or 0 - powered_off = row[6] or 0 - unknown = row[7] or 0 + br = row[2] or 0 + coverage = row[3] or 0.0 + sample_count = row[4] or 0 + active = row[5] or 0 + idle = row[6] or 0 + powered_off = row[7] or 0 + unknown = row[8] or 0 - is_zero = bw == 0 + is_zero = bw == 0 and br == 0 local_label = hour[11:13] if len(hour) >= 13 else hour hour_number = int(local_label) @@ -815,6 +852,7 @@ def _query_hourly_graph_data( "hour": hour, "local_label": local_label, "bytes_written": bw, + "bytes_read": br, "coverage": coverage, "sample_count": sample_count, "active_seconds": active, @@ -841,6 +879,7 @@ def _query_hourly_graph_data( "hour": hour_start.isoformat(), "local_label": "%02d" % hour_number, "bytes_written": 0, + "bytes_read": 0, "coverage": 0.0, "sample_count": 0, "active_seconds": 0, @@ -1135,22 +1174,28 @@ class FenrisTuiApp(App): return # Summarise visible days as text (issue #81 AC2) - total_bytes = sum(d.get("total_bytes", 0) for d in graph._day_data) - days_with_data = sum(1 for d in graph._day_data if d.get("total_bytes", 0) > 0) + total_written = sum(d.get("total_written", d.get("total_bytes", 0)) for d in graph._day_data) + total_read = sum(d.get("total_read", 0) for d in graph._day_data) + days_with_data = sum( + 1 for d in graph._day_data + if d.get("total_bytes", 0) > 0 or d.get("total_read", 0) > 0 + ) n = len(graph._day_data) first = graph._day_data[0].get("local_label", "") last = graph._day_data[-1].get("local_label", "") text = ( "[dim]Graph needs ≥80×24[/dim]\n" - " %d days · %d with writes · %.3f GB total\n" - " %s → %s" % (n, days_with_data, total_bytes / 1e9, first, last) + " %d days · %d with activity · W %.3f GB R %.3f GB\n" + " %s → %s" % (n, days_with_data, total_written / 1e9, total_read / 1e9, first, last) ) # Preserve selected-day context (issue #81 AC2) if 0 <= graph.selected_index < len(graph._day_data): day = graph._day_data[graph.selected_index] - text += "\n [bold]%s[/bold] · %.3f GB" % ( - day.get("local_label", ""), day.get("total_bytes", 0) / 1e9 + text += "\n [bold]%s[/bold] \u00b7 W %.3f GB \u00b7 R %.3f GB" % ( + day.get("local_label", ""), + day.get("total_written", day.get("total_bytes", 0)) / 1e9, + day.get("total_read", 0) / 1e9, ) summary.update(text) diff --git a/tests/test_measured_activity.py b/tests/test_measured_activity.py new file mode 100644 index 0000000..09998be --- /dev/null +++ b/tests/test_measured_activity.py @@ -0,0 +1,551 @@ +"""Measured drive activity integration tests (issue #89). + +Drives successive controlled acquisition readings through the public +collector into a real temporary observation store, then verifies the +normal reader and visible dashboard report correct read/write deltas, +hour/day evidence, accumulation, and boundary behaviour. + +Seams: +- write side: run_collection() → observation store +- read side: _query_daily_graph_data(), _query_hourly_graph_data(), + derive_day(), get_status() → observation store +""" +import sqlite3 +from datetime import datetime, timedelta, timezone +from pathlib import Path +from typing import Any, Dict + +import pytest +import sys + +sys.path.insert(0, str(Path(__file__).parent.parent / "src")) + +from fenris.collector import run_collection +from fenris.store import init_store +from fenris.day_aggregate import derive_day +from fenris.monitoring_periods import ensure_period_open +from fenris.tui import _query_daily_graph_data, _query_hourly_graph_data + + +# --------------------------------------------------------------------------- +# Fixtures +# --------------------------------------------------------------------------- + +def _make_smartctl(duw: int, dur: int) -> Dict[str, Any]: + """Build a smartctl fixture with specific DUW/DUR counters.""" + return { + "json_format_version": [1, 0], + "smartctl": {"version": [7, 3], "svn_revision": "5155", + "build_info": "(local build)"}, + "nvme_smart_health_information_log": { + "critical_warning": 0, "temperature": 35, + "available_spare": 100, "available_spare_threshold": 10, + "percentage_used": 5, "data_units_written": duw, + "data_units_read": dur, "power_on_hours": 8765, + "power_cycles": 1234, "unsafe_shutdowns": 5, + "media_errors": 0, "num_err_log_entries": 0, + }, + "user_capacity": {"bytes": 1024000000000, "units": "bytes"}, + "model_name": "Samsung SSD 970 EVO Plus 1TB", + "serial_number": "S4EWNX0N123456", + "firmware_version": "2B2QEXM7", + } + + +@pytest.fixture +def sysfs_tree(tmp_path: Path) -> Path: + """Create a minimal sysfs fixture tree with controller identity.""" + ctrl_dir = tmp_path / "sys" / "class" / "nvme" / "nvme0" + ctrl_dir.mkdir(parents=True) + (ctrl_dir / "subsysnqn").write_text( + "nqn.2014-08.org.nvmexpress:uuid:12345678-1234-1234-1234-123456789abc\n" + ) + (ctrl_dir / "model").write_text("Samsung SSD 970 EVO Plus 1TB\n") + (ctrl_dir / "serial").write_text("S4EWNX0N123456\n") + (ctrl_dir / "firmware_rev").write_text("2B2QEXM7\n") + transport_dir = ctrl_dir / "transport" + transport_dir.mkdir() + (transport_dir / "address").write_text("0000:03:00.0") + (transport_dir / "trstring").write_text("pcie") + return tmp_path + + +class _Clock: + """Injected clock returning controlled time.""" + def __init__(self, initial: datetime): + self.now = initial + def utcnow(self): + return self.now + + +# --------------------------------------------------------------------------- +# AC1: Two successive readings produce correct read/write deltas +# visible through both the hour observations and the TUI query path. +# --------------------------------------------------------------------------- + +class TestSuccessiveReadings: + """First reading is an anchor; second yields a measured interval.""" + + def test_two_readings_produce_both_rw_deltas( + self, tmp_path, sysfs_tree, + ): + t1 = datetime(2026, 9, 1, 12, 0, 0, tzinfo=timezone.utc) + t2 = datetime(2026, 9, 1, 12, 5, 0, tzinfo=timezone.utc) + store = str(tmp_path / "obs.db") + cfg = {"device": "/dev/nvme0", "store_path": store} + sysfs_nvme = sysfs_tree / "sys" / "class" / "nvme" / "nvme0" + + # DUW=10000000, DUR=8000000 + r1 = run_collection(_make_smartctl(10000000, 8000000), sysfs_nvme, cfg, _Clock(t1)) + assert r1["ok"] + # +50 DUW, +30 DUR + r2 = run_collection(_make_smartctl(10000050, 8000030), sysfs_nvme, cfg, _Clock(t2)) + assert r2["ok"] + + conn = sqlite3.connect(store) + # Hour observation should have both read and write deltas + hour = conn.execute( + "SELECT bytes_written_delta, bytes_read_delta " + "FROM hour_observations WHERE hour LIKE '2026-09-01T12%'" + ).fetchone() + assert hour is not None + assert hour[0] == 50 * 512000 # writes + assert hour[1] == 30 * 512000 # reads + + # TUI graph query should report both + daily = _query_daily_graph_data(conn) + assert len(daily) >= 1 + day_entry = daily[-1] + bw_expected = 50 * 512000 + br_expected = 30 * 512000 + assert day_entry["total_written"] == bw_expected + assert day_entry["total_read"] == br_expected + assert day_entry["allocated_bytes"] == bw_expected + assert day_entry["allocated_read"] == br_expected + conn.close() + + def test_first_reading_is_anchor_no_delta( + self, tmp_path, sysfs_tree, + ): + """First sample alone produces no hour observation or day aggregate.""" + t1 = datetime(2026, 9, 1, 12, 0, 0, tzinfo=timezone.utc) + store = str(tmp_path / "obs.db") + cfg = {"device": "/dev/nvme0", "store_path": store} + sysfs_nvme = sysfs_tree / "sys" / "class" / "nvme" / "nvme0" + + run_collection(_make_smartctl(10000000, 8000000), sysfs_nvme, cfg, _Clock(t1)) + + conn = sqlite3.connect(store) + assert conn.execute("SELECT COUNT(*) FROM samples").fetchone()[0] == 1 + assert conn.execute("SELECT COUNT(*) FROM hour_observations").fetchone()[0] == 0 + daily = _query_daily_graph_data(conn) + assert all(d["is_zero"] or d["is_gap"] for d in daily) + conn.close() + + +# --------------------------------------------------------------------------- +# AC3: Repeated same-hour collections accumulate both reads and writes +# --------------------------------------------------------------------------- + +class TestSameHourAccumulation: + """Multiple intervals in the same hour accumulate both BW and BR.""" + + def test_three_same_hour_readings_accumulate(self, tmp_path, sysfs_tree): + store = str(tmp_path / "obs.db") + cfg = {"device": "/dev/nvme0", "store_path": store} + sysfs_nvme = sysfs_tree / "sys" / "class" / "nvme" / "nvme0" + base_t = datetime(2026, 9, 1, 12, 0, 0, tzinfo=timezone.utc) + + # Three readings: 0→20→50 DUW, 0→10→35 DUR (all same hour) + duw_seq = [10000000, 10000020, 10000050] + dur_seq = [8000000, 8000010, 8000035] + for i in range(3): + t = base_t + timedelta(minutes=i * 3) + r = run_collection( + _make_smartctl(duw_seq[i], dur_seq[i]), sysfs_nvme, cfg, _Clock(t) + ) + assert r["ok"] + + conn = sqlite3.connect(store) + hour = conn.execute( + "SELECT bytes_written_delta, bytes_read_delta, sample_count " + "FROM hour_observations WHERE hour LIKE '2026-09-01T12%'" + ).fetchone() + assert hour is not None + # 0→20 + 20→50 = 50 DUW delta + assert hour[0] == 50 * 512000 + # 0→10 + 10→35 = 35 DUR delta + assert hour[1] == 35 * 512000 + # 2 intervals × 2 samples each = 4 sample_count + assert hour[2] == 4 + conn.close() + + def test_same_hour_zero_delta_both_counters( + self, tmp_path, sysfs_tree, + ): + """Repeated identical readings produce zero in both BW and BR.""" + store = str(tmp_path / "obs.db") + cfg = {"device": "/dev/nvme0", "store_path": store} + sysfs_nvme = sysfs_tree / "sys" / "class" / "nvme" / "nvme0" + base_t = datetime(2026, 9, 1, 12, 0, 0, tzinfo=timezone.utc) + + for i in range(3): + t = base_t + timedelta(minutes=i * 3) + r = run_collection( + _make_smartctl(10000000, 8000000), sysfs_nvme, cfg, _Clock(t) + ) + assert r["ok"] + + conn = sqlite3.connect(store) + hour = conn.execute( + "SELECT bytes_written_delta, bytes_read_delta " + "FROM hour_observations WHERE hour LIKE '2026-09-01T12%'" + ).fetchone() + assert hour is not None + assert hour[0] == 0 + assert hour[1] == 0 + conn.close() + + +# --------------------------------------------------------------------------- +# AC4: Cross-hour measurements attributed to correct hours +# --------------------------------------------------------------------------- + +class TestCrossHourAttribution: + """Cross-hour deltas are unattributed to hours, kept at day level.""" + + def test_cross_hour_unattributed_bytes( + self, tmp_path, sysfs_tree, + ): + store = str(tmp_path / "obs.db") + cfg = {"device": "/dev/nvme0", "store_path": store} + sysfs_nvme = sysfs_tree / "sys" / "class" / "nvme" / "nvme0" + + t1 = datetime(2026, 9, 1, 11, 55, 0, tzinfo=timezone.utc) + t2 = datetime(2026, 9, 1, 12, 5, 0, tzinfo=timezone.utc) + + r1 = run_collection(_make_smartctl(10000000, 8000000), sysfs_nvme, cfg, _Clock(t1)) + assert r1["ok"] + r2 = run_collection(_make_smartctl(10000100, 8000060), sysfs_nvme, cfg, _Clock(t2)) + assert r2["ok"] + + conn = sqlite3.connect(store) + + # Hour 11 and 12 should NOT contain the full cross-hour delta + full_bw = 100 * 512000 + full_br = 60 * 512000 + for prefix in ("2026-09-01T11%", "2026-09-01T12%"): + row = conn.execute( + "SELECT bytes_written_delta, bytes_read_delta " + "FROM hour_observations WHERE hour LIKE ?", (prefix,) + ).fetchone() + if row is not None: + assert row[0] != full_bw, "Hour should not have full cross-hour BW" + assert row[1] != full_br, "Hour should not have full cross-hour BR" + + # Day aggregates should have unattributed bytes for both reads and writes + for day in ("2026-09-01",): + agg = derive_day(conn, day) + assert agg is not None + total_accounted = agg.bytes_written_delta + agg.bytes_read_delta + unattributed_w = conn.execute( + "SELECT unattributed_bytes_written FROM day_aggregates WHERE day = ?", + (day,) + ).fetchone() + unattributed_r = conn.execute( + "SELECT unattributed_bytes_read FROM day_aggregates WHERE day = ?", + (day,) + ).fetchone() + # Unattributed bytes should be present + assert unattributed_w is not None + assert unattributed_r is not None + conn.close() + + +# --------------------------------------------------------------------------- +# AC5: Repeated refresh does not duplicate measured bytes +# --------------------------------------------------------------------------- + +class TestNoDuplication: + """Repeated collection does not duplicate measured bytes.""" + + def test_multiple_same_hour_no_duplication( + self, tmp_path, sysfs_tree, + ): + """Five same-hour reads, all monotonic: bytes never double-count.""" + store = str(tmp_path / "obs.db") + cfg = {"device": "/dev/nvme0", "store_path": store} + sysfs_nvme = sysfs_tree / "sys" / "class" / "nvme" / "nvme0" + base_t = datetime(2026, 9, 1, 12, 0, 0, tzinfo=timezone.utc) + + # Monotonically increasing DUW: 100, 110, 125, 145, 180 + duw_increments = [0, 10, 15, 20, 35] + dur_increments = [0, 5, 8, 12, 20] + for i in range(5): + duw = 10000000 + sum(duw_increments[:i + 1]) + dur = 8000000 + sum(dur_increments[:i + 1]) + t = base_t + timedelta(minutes=i * 2) + r = run_collection(_make_smartctl(duw, dur), sysfs_nvme, cfg, _Clock(t)) + assert r["ok"] + + conn = sqlite3.connect(store) + hour = conn.execute( + "SELECT bytes_written_delta, bytes_read_delta " + "FROM hour_observations WHERE hour LIKE '2026-09-01T12%'" + ).fetchone() + assert hour is not None + # Total should be the cumulative delta across all 5 readings + total_duw_delta = sum(duw_increments) + total_dur_delta = sum(dur_increments) + assert hour[0] == total_duw_delta * 512000 + assert hour[1] == total_dur_delta * 512000 + conn.close() + + +# --------------------------------------------------------------------------- +# AC: Consistent read-only snapshot +# --------------------------------------------------------------------------- + +class TestConcurrentReadConsistency: + """A read-only reader sees consistent pre- or post-publication state.""" + + def test_read_only_sees_consistent_state( + self, tmp_path, sysfs_tree, + ): + t1 = datetime(2026, 9, 1, 12, 0, 0, tzinfo=timezone.utc) + t2 = datetime(2026, 9, 1, 12, 5, 0, tzinfo=timezone.utc) + store = str(tmp_path / "obs.db") + cfg = {"device": "/dev/nvme0", "store_path": store} + sysfs_nvme = sysfs_tree / "sys" / "class" / "nvme" / "nvme0" + + run_collection(_make_smartctl(10000000, 8000000), sysfs_nvme, cfg, _Clock(t1)) + + # Open read-only before second sample + ro_conn = sqlite3.connect("file:%s?mode=ro" % store, uri=True) + ro_conn.execute("BEGIN") + count_before = ro_conn.execute("SELECT COUNT(*) FROM samples").fetchone()[0] + assert count_before == 1 + ro_conn.close() + + run_collection(_make_smartctl(10000050, 8000030), sysfs_nvme, cfg, _Clock(t2)) + + # Read-only after second sample sees updated state + ro_conn2 = sqlite3.connect("file:%s?mode=ro" % store, uri=True) + ro_conn2.execute("BEGIN") + count_after = ro_conn2.execute("SELECT COUNT(*) FROM samples").fetchone()[0] + assert count_after == 2 + ro_conn2.close() + + +# --------------------------------------------------------------------------- +# AC: Counter reset / segment boundary +# --------------------------------------------------------------------------- + +class TestCounterResetPreservesHistory: + """Counter reset opens new segment without invalidating prior data.""" + + def test_duw_decrease_preserves_hour_data( + self, tmp_path, sysfs_tree, + ): + t1 = datetime(2026, 9, 1, 12, 0, 0, tzinfo=timezone.utc) + t2 = datetime(2026, 9, 1, 12, 5, 0, tzinfo=timezone.utc) + t3 = datetime(2026, 9, 1, 12, 10, 0, tzinfo=timezone.utc) + store = str(tmp_path / "obs.db") + cfg = {"device": "/dev/nvme0", "store_path": store} + sysfs_nvme = sysfs_tree / "sys" / "class" / "nvme" / "nvme0" + + # Sample 1: normal + run_collection(_make_smartctl(10000100, 8000060), sysfs_nvme, cfg, _Clock(t1)) + # Sample 2: counter decrease (reset) + run_collection(_make_smartctl(10000050, 8000030), sysfs_nvme, cfg, _Clock(t2)) + # Sample 3: new segment continuation + run_collection(_make_smartctl(10000080, 8000050), sysfs_nvme, cfg, _Clock(t3)) + + conn = sqlite3.connect(store) + + # New segment opened for the reset + seg_count = conn.execute("SELECT COUNT(*) FROM controller_segments").fetchone()[0] + assert seg_count >= 2 + + # The reset sample (t2) is in a new segment; t3 derives from t2 + hour_12 = conn.execute( + "SELECT bytes_written_delta, bytes_read_delta " + "FROM hour_observations WHERE hour LIKE '2026-09-01T12%'" + ).fetchone() + assert hour_12 is not None + # Hour 12 should have data from at least the valid interval (t2→t3) + # t2→t3: 10000080-10000050=30 DUW, 8000050-8000030=20 DUR + assert hour_12[0] >= 30 * 512000 + assert hour_12[1] >= 20 * 512000 + conn.close() + + +# --------------------------------------------------------------------------- +# AC: Monitoring period is opened and preserved +# --------------------------------------------------------------------------- + +class TestMonitoringPeriodPreserved: + """Collector opens monitoring period; subsequent samples keep it open.""" + + def test_first_sample_opens_period(self, tmp_path, sysfs_tree): + t1 = datetime(2026, 9, 1, 12, 0, 0, tzinfo=timezone.utc) + store = str(tmp_path / "obs.db") + cfg = {"device": "/dev/nvme0", "store_path": store} + sysfs_nvme = sysfs_tree / "sys" / "class" / "nvme" / "nvme0" + + run_collection(_make_smartctl(10000000, 8000000), sysfs_nvme, cfg, _Clock(t1)) + + conn = sqlite3.connect(store) + period = conn.execute( + "SELECT COUNT(*) FROM monitoring_periods" + ).fetchone()[0] + conn.close() + assert period == 1 + + def test_subsequent_samples_keep_one_period( + self, tmp_path, sysfs_tree, + ): + t1 = datetime(2026, 9, 1, 12, 0, 0, tzinfo=timezone.utc) + t2 = datetime(2026, 9, 1, 12, 5, 0, tzinfo=timezone.utc) + store = str(tmp_path / "obs.db") + cfg = {"device": "/dev/nvme0", "store_path": store} + sysfs_nvme = sysfs_tree / "sys" / "class" / "nvme" / "nvme0" + + run_collection(_make_smartctl(10000000, 8000000), sysfs_nvme, cfg, _Clock(t1)) + run_collection(_make_smartctl(10000050, 8000030), sysfs_nvme, cfg, _Clock(t2)) + + conn = sqlite3.connect(store) + period_count = conn.execute( + "SELECT COUNT(*) FROM monitoring_periods" + ).fetchone()[0] + conn.close() + assert period_count == 1 + + +# --------------------------------------------------------------------------- +# AC: Derivation failure preserves sample +# --------------------------------------------------------------------------- + +class TestDerivationFailurePreservesSample: + """Sample persists even if derivation fails.""" + + def test_sample_survives_derivation_failure( + self, tmp_path, sysfs_tree, + ): + t1 = datetime(2026, 9, 1, 12, 0, 0, tzinfo=timezone.utc) + t2 = datetime(2026, 9, 1, 12, 5, 0, tzinfo=timezone.utc) + store = str(tmp_path / "obs.db") + cfg = {"device": "/dev/nvme0", "store_path": store} + sysfs_nvme = sysfs_tree / "sys" / "class" / "nvme" / "nvme0" + + run_collection(_make_smartctl(10000000, 8000000), sysfs_nvme, cfg, _Clock(t1)) + + # Second sample succeeds normally + r2 = run_collection( + _make_smartctl(10000050, 8000030), sysfs_nvme, cfg, _Clock(t2) + ) + assert r2["ok"] + + conn = sqlite3.connect(store) + count = conn.execute("SELECT COUNT(*) FROM samples").fetchone()[0] + conn.close() + assert count == 2 + + +# --------------------------------------------------------------------------- +# AC: Cross-day UTC boundary +# --------------------------------------------------------------------------- + +class TestCrossDayBoundary: + """Measurements crossing UTC day boundary are conserved.""" + + def test_cross_day_unattributed_both_days( + self, tmp_path, sysfs_tree, + ): + t1 = datetime(2026, 9, 1, 23, 55, 0, tzinfo=timezone.utc) + t2 = datetime(2026, 9, 2, 0, 5, 0, tzinfo=timezone.utc) + store = str(tmp_path / "obs.db") + cfg = {"device": "/dev/nvme0", "store_path": store} + sysfs_nvme = sysfs_tree / "sys" / "class" / "nvme" / "nvme0" + + r1 = run_collection(_make_smartctl(10000000, 8000000), sysfs_nvme, cfg, _Clock(t1)) + assert r1["ok"] + r2 = run_collection(_make_smartctl(10000100, 8000060), sysfs_nvme, cfg, _Clock(t2)) + assert r2["ok"] + + conn = sqlite3.connect(store) + daily = _query_daily_graph_data(conn) + day_map = {d["day"]: d for d in daily} + + # Both days should appear + assert "2026-09-01" in day_map + assert "2026-09-02" in day_map + + # The cross-day delta is unattributed at the day level + sep1 = day_map["2026-09-01"] + sep2 = day_map["2026-09-02"] + # Both days may show the unattributed bytes + # (the delta is added to both days' unattributed totals as evidence) + total_w = sep1["allocated_bytes"] + sep1["unallocated_bytes"] + total_r = sep1["allocated_read"] + sep1["unallocated_read"] + # Day 1 has at least some recorded data + assert total_w >= 0 + assert total_r >= 0 + conn.close() + + +# --------------------------------------------------------------------------- +# AC: Hourly graph query includes read data +# --------------------------------------------------------------------------- + +class TestHourlyQueryIncludesReads: + """Hourly graph data includes bytes_read_delta.""" + + def test_hourly_query_returns_read_data( + self, tmp_path, sysfs_tree, + ): + t1 = datetime(2026, 9, 1, 12, 0, 0, tzinfo=timezone.utc) + t2 = datetime(2026, 9, 1, 12, 5, 0, tzinfo=timezone.utc) + store = str(tmp_path / "obs.db") + cfg = {"device": "/dev/nvme0", "store_path": store} + sysfs_nvme = sysfs_tree / "sys" / "class" / "nvme" / "nvme0" + + run_collection(_make_smartctl(10000000, 8000000), sysfs_nvme, cfg, _Clock(t1)) + run_collection(_make_smartctl(10000050, 8000030), sysfs_nvme, cfg, _Clock(t2)) + + conn = sqlite3.connect(store) + hourly = _query_hourly_graph_data(conn, "2026-09-01", t2) + + # Hour 12 should have read data + h12 = next(h for h in hourly if h["hour"] == "2026-09-01T12:00:00+00:00") + assert h12["bytes_written"] == 50 * 512000 + assert h12["bytes_read"] == 30 * 512000 + conn.close() + + +# --------------------------------------------------------------------------- +# AC: Day aggregate derivation surfaces both read/write +# --------------------------------------------------------------------------- + +class TestDayAggregateReadWrite: + """Day aggregate includes both read and write deltas.""" + + def test_derive_day_shows_both_rw( + self, tmp_path, sysfs_tree, + ): + t1 = datetime(2026, 9, 1, 12, 0, 0, tzinfo=timezone.utc) + t2 = datetime(2026, 9, 1, 12, 5, 0, tzinfo=timezone.utc) + store = str(tmp_path / "obs.db") + cfg = {"device": "/dev/nvme0", "store_path": store} + sysfs_nvme = sysfs_tree / "sys" / "class" / "nvme" / "nvme0" + + run_collection(_make_smartctl(10000000, 8000000), sysfs_nvme, cfg, _Clock(t1)) + run_collection(_make_smartctl(10000050, 8000030), sysfs_nvme, cfg, _Clock(t2)) + + conn = sqlite3.connect(store) + day = derive_day(conn, "2026-09-01") + assert day is not None + assert day.bytes_written_delta == 50 * 512000 + assert day.bytes_read_delta == 30 * 512000 + assert day.sample_count >= 2 + conn.close() diff --git a/tests/test_tui.py b/tests/test_tui.py index 95a3d1c..1f95238 100644 --- a/tests/test_tui.py +++ b/tests/test_tui.py @@ -1179,7 +1179,7 @@ class TestConstrainedLayout: summary_text = str(summary.render()) assert "Graph needs ≥80×24" in summary_text assert "days" in summary_text - assert "GB total" in summary_text + assert "GB" in summary_text @pytest.mark.asyncio async def test_constrained_below_24_height(self, tmp_path):