Compare commits

..
Author SHA1 Message Date
xavierkandCommandCodeBot 4f2f30abac feat(#76): anchor scenario windows at evidence endpoint T
Headline and scenario rates now describe exact evidence-supported
monitored spans anchored at the latest published usage-evidence
endpoint T, not clock_now. Reader refresh alone never moves T or
dilutes rates.

- Add horizon_reasons field to ScenarioRange for specific unavailability facts
- Modify _compute_horizon_rate to use exact trailing 7/28/90×86400-second starts from T
- Show specific reasons for affected horizons (e.g., "starts before earliest data")
- Update TUI and CLI to display horizon-specific reasons
- Add 6 new tests for evidence-anchored projection rates

Co-authored-by: CommandCodeBot <noreply@commandcode.ai>
2026-09-14 11:11:50 +05:30
xavierkandCommandCodeBot 5d916ee97f feat: add interactive daily writes bar graph with hourly drill-down (issue #75)
Replace the static sparkline with an interactive block-glyph bar graph
that supports writes-only daily bars, range switching (7/14/28/90 days),
day selection, and hourly drill-down.  No plotting dependency.

Co-authored-by: CommandCodeBot <noreply@commandcode.ai>
2026-09-14 10:47:37 +05:30
xavierk 6917a658cb feat: implement repair and retention for observation history (issue #74) 2026-09-14 03:54:46 +05:30
xavierk d790ff84c5 feat(#73): publish trustworthy first usage history
- Schema migration 1-2: add segment_id to samples, unattributed bytes to day_aggregates

- Collector now derives hour observations and day aggregates from sample pairs

- Cross-hour deltas tracked as unattributed (no proportional allocation)

- Display states: 0 samples -> awaiting first, 1 sample -> awaiting another

- Monitoring period ensured open on each collection run

- Derivation failures preserve prior history

Closes #73
2026-09-14 03:19:08 +05:30
26 changed files with 3318 additions and 2134 deletions
-45
View File
@@ -1,45 +0,0 @@
# PROTOTYPE — graph encoding & hourly drill-down (throwaway)
Answers wayfinder ticket **Choose daily graph encoding and hourly drill-down**
on map **Fenris TUI polish and hourly history**.
Not production code. Do not merge onto main.
## Question
Which graph encoding best communicates hourly writes and a day's usage
without misleading: daily bars vs one candle per day vs rolling hourly strip,
with range -> selected day -> hourly detail -> back navigation.
## Run
./run # or: python3 tui_prototype.py
## Variants (switch with [ and ])
- **A — Daily bars**: one bar per local display day (bytes written), drill
into 24 hourly bars with Enter / click.
- **B — Daily candles**: one candle per day derived from hourly values
(open = first evidenced hour, close = last, high/low = max/min hourly).
Same drill-down.
- **C — Rolling hourly strip**: continuous last-72-hours columns with day
separators; click an hour for its readout (no day level).
## Data states encoded (synthetic)
measured zero (0 B) · gap / no evidence · unallocated usage (measured,
attribution unknown) · partial day ("so far") · deliberate-disable pause band ·
day before monitoring began.
Gaps never render as zero. Unallocated is a separate visual segment, never
spread into hours.
## Framework facts verified
Textual 8.2.8 (installed): no BarChart widget (roadmap-only), Sparkline is
non-interactive. Bars/candles here are a custom Static renderable; mouse
clicks map via widget-local x. No new dependencies.
## Screenshots
screenshots/ holds SVG exports at 80x24 and 140x40 per variant.
Regenerate with: python3 capture.py
-31
View File
@@ -1,31 +0,0 @@
#!/usr/bin/env python3
# PROTOTYPE helper — regenerate screenshots/*.svg via headless Textual run.
import asyncio
from tui_prototype import GraphPrototype
async def shot(app, size, presses, path):
async with app.run_test(size=size) as pilot:
for p in presses:
await pilot.press(p)
svg = app.export_screenshot()
with open(path, "w") as fh:
fh.write(svg)
print("wrote", path)
async def main():
for w, h in ((80, 24), (140, 40)):
for v in "ABC":
presses = ["]"] * "ABC".index(v)
await shot(GraphPrototype(), (w, h), presses,
"screenshots/v%s_%dx%d_range.svg" % (v, w, h))
# drill-down detail (variant A) and a mid-range selection (variant B)
await shot(GraphPrototype(), (80, 24), ["left", "enter"],
"screenshots/vA_80x24_hourly-detail.svg")
await shot(GraphPrototype(), (140, 40), ["]", "left", "left"],
"screenshots/vB_140x40_selected-gap-day.svg")
if __name__ == "__main__":
asyncio.run(main())
-3
View File
@@ -1,3 +0,0 @@
#!/bin/sh
# PROTOTYPE launcher
exec python3 "$(dirname "$0")/tui_prototype.py" "$@"
File diff suppressed because one or more lines are too long

Before

Width:  |  Height:  |  Size: 39 KiB

File diff suppressed because one or more lines are too long

Before

Width:  |  Height:  |  Size: 46 KiB

File diff suppressed because one or more lines are too long

Before

Width:  |  Height:  |  Size: 35 KiB

File diff suppressed because one or more lines are too long

Before

Width:  |  Height:  |  Size: 39 KiB

File diff suppressed because one or more lines are too long

Before

Width:  |  Height:  |  Size: 39 KiB

File diff suppressed because one or more lines are too long

Before

Width:  |  Height:  |  Size: 35 KiB

File diff suppressed because one or more lines are too long

Before

Width:  |  Height:  |  Size: 84 KiB

File diff suppressed because one or more lines are too long

Before

Width:  |  Height:  |  Size: 78 KiB

-488
View File
@@ -1,488 +0,0 @@
#!/usr/bin/env python3
# PROTOTYPE — throwaway. Wayfinder ticket "Choose daily graph encoding and
# hourly drill-down" on map "Fenris TUI polish and hourly history".
# Synthetic data only. Not production code; no tests by design.
#
# Variants (switch with [ / ]):
# A daily bars (day -> Enter -> hourly bars -> Esc back)
# B daily candles (OHLC derived from hourly values, same drill-down)
# C rolling 72-hour strip with day separators (no day level)
#
# Encoding contract under test (accepted hourly-history decision):
# measured zero = 0 B, distinct from gap (no evidence), distinct from
# unallocated usage (measured total, hour attribution unknown).
from __future__ import annotations
import time
from dataclasses import dataclass, field
from datetime import datetime, timedelta
from textual import events
from textual.app import App, ComposeResult
from textual.binding import Binding
from textual.containers import Vertical
from textual.widgets import Static
# ---------------------------------------------------------------------------
# Synthetic evidence
# ---------------------------------------------------------------------------
@dataclass
class Hour:
gb: float = 0.0 # measured bytes written; 0.0 == measured zero
coverage: float = 1.0 # share of hour classified known
@dataclass
class Day:
hours: list = field(default_factory=list) # list[Hour | None]; None = gap
unallocated: float = 0.0 # GB measured, attribution unknown
elapsed_h: float = 24.0 # elapsed span (partial < 24)
note: str = ""
pause: str = "" # deliberate-disable band
@property
def evidenced(self):
return [h for h in self.hours if h is not None]
@property
def allocated_gb(self):
return sum(h.gb for h in self.evidenced)
@property
def gap_hours(self):
return sum(1 for h in self.hours if h is None)
def _profile(mult=1.0, zero=False):
shape = [0.3, 0.2, 0.1, 0.1, 0.2, 0.5, 1.4, 2.6, 3.1, 2.2,
1.8, 2.0, 2.4, 3.0, 3.4, 2.8, 2.1, 1.9, 2.3, 1.6,
1.1, 0.8, 0.5, 0.4]
return [Hour(0.0 if zero else round(v * mult, 2), 1.0) for v in shape]
def synth_days(now):
"""14 local display days ending today (partial). Every ticket state present."""
days = []
days.append(Day([None] * 24, note="no evidence — before monitoring began"))
days.append(Day(_profile(1.0)))
g = _profile(0.9)
for i in range(10, 16):
g[i] = None
days.append(Day(g, note="gap 10:00–16:00 (no evidence, not zero)"))
days.append(Day(_profile(0.0), note="measured zero day"))
days.append(Day(_profile(1.2), unallocated=8.5,
note="8.5 GB unallocated (interval attribution unknown)"))
days.append(Day(_profile(4.0), note="heavy write day"))
days.append(Day(_profile(0.7),
pause="paused 09:00–11:00 (deliberate disable — excluded)"))
days.append(Day(_profile(1.1)))
days.append(Day(_profile(0.9)))
m = _profile(1.0)
m[3] = Hour(0.0, 1.0)
m[17] = None
days.append(Day(m, note="zero hour 03:00, gap hour 17:00"))
days.append(Day(_profile(1.3), unallocated=1.9))
days.append(Day(_profile(1.0)))
days.append(Day(_profile(0.8)))
days.append(Day(_profile(1.5)[:14], elapsed_h=14,
note="today — 14 h elapsed, so far"))
assert len(days) == 14
return days
# ---------------------------------------------------------------------------
# Rendering (Rich markup into Static)
# ---------------------------------------------------------------------------
AMBER = "#e0a458" # default-preset graph colour (approved input)
UNALLOC = "#8a7f70" # measured, attribution unknown
GAP = "#5f6b73"
ZERO = "#9bb0bf"
SELECT = "#f2d5a0"
FILL = "█"
EIGHTHS = " ▁▂▃▄▅▆▇█"
def _fmt_gb(v):
if v >= 100:
return "%d GB" % v
return "%s GB" % ("%.1f" % v if v < 10 else "%.0f" % v)
def render_days_bars(days, height, sel):
"""Variant A: one bar per day. Unallocated stacked on top, visually separate."""
n = len(days)
mx = max((d.allocated_gb + d.unallocated) for d in days) or 1.0
grid = [[" "] * n for _ in range(height)]
style = [[None] * n for _ in range(height)]
for i, d in enumerate(days):
if not d.evidenced and not d.unallocated:
for r in range(height):
grid[r][i], style[r][i] = "·", GAP
continue
rows = []
ua_cells = d.unallocated / mx * height
for _ in range(int(ua_cells + 1e-9)):
rows.append(("▒", UNALLOC))
if ua_cells - int(ua_cells) > 0.12:
rows.append((EIGHTHS[int((ua_cells % 1) * 8)], UNALLOC))
cells = (d.allocated_gb + d.unallocated) / mx * height
for _ in range(int(cells + 1e-9)):
rows.append((FILL, AMBER))
if cells - int(cells) > 0.12:
rows.append((EIGHTHS[int((cells % 1) * 8)], AMBER))
if not rows:
rows = [("·", ZERO)]
rr = height - 1
for ch, st in rows:
if rr < 0:
break
grid[rr][i], style[rr][i] = ch, st
rr -= 1
if d.elapsed_h < 24:
rr = height - len(rows)
if 0 <= rr < height:
grid[rr][i], style[rr][i] = "┄", SELECT
lines = []
for r in range(height):
parts = []
for i in range(n):
ch, st = grid[r][i], style[r][i]
parts.append(ch if st is None else "[%s]%s[/]" % (st, ch))
parts.append(" ")
lines.append("".join(parts).rstrip())
marks = [" "] * (2 * n)
marks[2 * sel] = "▼"
lines.insert(0, "".join(marks).rstrip())
return lines
def render_days_candles(days, height, sel):
"""Variant B: one candle per day. open = first evidenced hour, close = last,
high/low = max/min hourly GB. Body solid when close >= open, hollow otherwise;
partial day renders dashed. 3-wide slot + 1 space."""
n = len(days)
stats = []
for d in days:
ev = [h.gb for h in d.evidenced]
stats.append(None if not ev else
dict(open=ev[0], close=ev[-1], hi=max(ev), lo=min(ev)))
mx = max((s["hi"] for s in stats if s), default=1.0) or 1.0
def row(v):
return height - 1 - min(int(v / mx * (height - 1)), height - 1)
cols = [] # per day: list[height] of 3-char cell strings + styles
for i, s in enumerate(stats):
d = days[i]
if s is None:
cols.append([("╎╎╎", GAP)] * height)
continue
top, bot = row(s["hi"]), row(s["lo"])
body_hi, body_lo = row(max(s["open"], s["close"])), row(min(s["open"], s["close"]))
solid = s["close"] >= s["open"]
partial = d.elapsed_h < 24
body_fill = "▒" if partial else (FILL if solid else "░")
body_st = AMBER if (solid or partial) else GAP
col = []
for r in range(height):
if body_lo <= r <= body_hi:
col.append((body_fill * 3, body_st))
elif top <= r <= bot:
col.append((" │ ", AMBER if solid else GAP))
else:
col.append((" ", None))
cols.append(col)
lines = [(" " * (4 * sel)) + "▼"]
for r in range(height):
parts = []
for i in range(n):
cells, st = cols[i][r]
parts.append(cells if st is None else "[%s]%s[/]" % (st, cells))
parts.append(" ")
lines.append("".join(parts).rstrip())
return lines
def render_hours_bars(day, height, sel):
"""Day detail: one bar per hour. Gaps never zero."""
n = len(day.hours)
mx = max([h.gb for h in day.evidenced] + [day.unallocated, 1e-9])
grid = [[" "] * n for _ in range(height)]
style = [[None] * n for _ in range(height)]
for i, h in enumerate(day.hours):
if h is None:
grid[height - 1][i], style[height - 1][i] = "░", GAP
continue
if h.gb == 0:
grid[height - 1][i], style[height - 1][i] = "·", ZERO
continue
cells = h.gb / mx * height
rr = height - 1
for _ in range(int(cells + 1e-9)):
if rr < 0:
break
grid[rr][i], style[rr][i] = FILL, AMBER
rr -= 1
if cells % 1 > 0.12 and rr >= 0:
grid[rr][i], style[rr][i] = EIGHTHS[int((cells % 1) * 8)], AMBER
lines = []
for r in range(height):
parts = []
for i in range(n):
ch, st = grid[r][i], style[r][i]
parts.append(ch if st is None else "[%s]%s[/]" % (st, ch))
lines.append("".join(parts).rstrip())
marks = [" "] * n
marks[sel % n] = "▼"
lines.insert(0, "".join(marks).rstrip())
return lines
def render_strip(days, height, sel):
"""Variant C: rolling hourly strip, last 72 h, day separators."""
flat = [] # (day_idx, hour_idx, Hour|None|'future')
for di, d in enumerate(days[-4:]):
real = len(days) - 4 + di
for hi in range(24):
h = d.hours[hi] if hi < len(d.hours) else "future"
flat.append((real, hi, h))
flat = flat[-72:]
mx = max([f[2].gb for f in flat if f[2] not in (None, "future")] + [1e-9])
cols = []
for _, _, h in flat:
col = [(" ", None)] * height
if h == "future":
col[height - 1] = ("·", None)
elif h is None:
col[height - 1] = ("░", GAP)
elif h.gb == 0:
col[height - 1] = ("·", ZERO)
else:
cells = h.gb / mx * height
rr = height - 1
for _ in range(int(cells + 1e-9)):
if rr < 0:
break
col[rr] = (FILL, AMBER)
rr -= 1
if cells % 1 > 0.12 and rr >= 0:
col[rr] = (EIGHTHS[int((cells % 1) * 8)], AMBER)
cols.append(col)
lines = []
for r in range(height):
parts = []
prev = None
for i, f in enumerate(flat):
if prev is not None and f[0] != prev:
parts.append("[%s]│[/]" % SELECT)
prev = f[0]
ch, st = cols[i][r]
parts.append(ch if st is None else "[%s]%s[/]" % (st, ch))
lines.append("".join(parts).rstrip())
marks = []
prev = None
for i, f in enumerate(flat):
if prev is not None and f[0] != prev:
marks.append(" ")
prev = f[0]
marks.append("▼" if i == sel else " ")
lines.insert(0, "".join(marks).rstrip())
return lines, flat
# ---------------------------------------------------------------------------
# App
# ---------------------------------------------------------------------------
VARIANT_NAMES = {"A": "Daily bars", "B": "Daily candles", "C": "Rolling 72 h strip"}
GH = 10 # graph rows
class GraphPrototype(App):
TITLE = "PROTOTYPE — usage-history graph encoding"
CSS = """
#col { height: 100%; }
#hdr { height: 2; color: $text-muted; }
#graph { height: auto; padding: 0 1; }
#readout { height: auto; padding: 0 1; }
#legend { height: 3; padding: 0 1; color: $text-muted; }
#footer { dock: bottom; height: 2; background: $panel; padding: 0 1; }
"""
BINDINGS = [
Binding("[", "prev_variant", "variant ←"),
Binding("]", "next_variant", "variant →"),
Binding("left", "left", "←"),
Binding("right", "right", "→"),
Binding("enter", "drill", "open"),
Binding("escape,backspace", "back", "back"),
Binding("q", "quit", "quit"),
]
def __init__(self):
super().__init__()
self.now = datetime.now().astimezone()
self.days = synth_days(self.now)
self.variant = "A"
self.view = "days" # days | hours (A/B); strip (C)
self.sel_day = len(self.days) - 1
self.sel_hour = 13
self.sel_strip = 71
def day_label(self, i):
d = (self.now - timedelta(days=len(self.days) - 1 - i)).date()
return d, d.strftime("%a %d %b")
def tz_label(self):
off = self.now.strftime("%z")
return "Local · UTC%s · %s" % (off[:3] + ":" + off[3:], time.tzname[0])
def day_readout(self, d, i):
_, lab = self.day_label(i)
if not d.evidenced and not d.unallocated:
return "%s — no evidence (not zero): %s" % (lab, d.note or "")
bits = ["%s — %s written" % (lab, _fmt_gb(d.allocated_gb))]
span = 24 if d.elapsed_h == 24 else len(d.hours)
bits.append("%d/%d h evidenced" % (span - d.gap_hours, span))
if d.unallocated:
bits.append("+ %s unallocated" % _fmt_gb(d.unallocated))
if d.elapsed_h < 24:
bits.append("partial · %d h elapsed · so far" % int(d.elapsed_h))
cov = sum(h.coverage for h in d.evidenced) / max(len(d.hours), 1)
bits.append("coverage %.0f%%" % (100 * cov))
txt = " · ".join(bits)
if d.note:
txt += " — %s" % d.note
if d.pause:
txt += " — %s" % d.pause
return txt
def compose(self) -> ComposeResult:
with Vertical(id="col"):
yield Static("", id="hdr")
yield Static("", id="graph")
yield Static("", id="readout")
yield Static("", id="legend")
yield Static("", id="footer")
def on_mount(self) -> None:
self.refresh_all()
def refresh_all(self):
self.query_one("#hdr", Static).update(
"🐺 Fenris · usage history · %s · writes per hour" % self.tz_label())
g = self.query_one("#graph", Static)
ro = self.query_one("#readout", Static)
if self.view == "hours":
d = self.days[self.sel_day]
lines = render_hours_bars(d, GH, self.sel_hour)
lines.append("".join(("%02d" % i) if i % 3 == 0 else " " for i in range(len(d.hours))))
lines.append(" midnight → 23:00 local")
g.update("\n".join(lines))
h = d.hours[self.sel_hour]
if h is None:
dtxt = "hour %02d:00 — no evidence (not zero)" % self.sel_hour
elif h.gb == 0:
dtxt = "hour %02d:00 — 0 B measured" % self.sel_hour
else:
dtxt = "hour %02d:00 — %s written · coverage %.0f%%" % (
self.sel_hour, _fmt_gb(h.gb), 100 * h.coverage)
_, lab = self.day_label(self.sel_day)
extra = ""
if d.unallocated:
extra = " · %s unallocated (shown separately, never spread into hours)" % _fmt_gb(d.unallocated)
ro.update("%s ▸ %s%s" % (lab, dtxt, extra))
elif self.variant in ("A", "B"):
lines = (render_days_bars(self.days, GH, self.sel_day) if self.variant == "A"
else render_days_candles(self.days, GH, self.sel_day))
_, last = self.day_label(len(self.days) - 1)
lines.append(" " + " ".join(
self.day_label(i)[1].split()[1] for i in range(len(self.days))))
g.update("\n".join(lines))
ro.update(self.day_readout(self.days[self.sel_day], self.sel_day))
else:
lines, flat = render_strip(self.days, GH, self.sel_strip)
g.update("\n".join(lines))
txt = ""
if 0 <= self.sel_strip < len(flat):
di, hi, h = flat[self.sel_strip]
_, lab = self.day_label(di)
if h == "future":
txt = "future — not counted"
elif h is None:
txt = "no evidence (not zero)"
elif h.gb == 0:
txt = "0 B measured"
else:
txt = "%s written · coverage %.0f%%" % (_fmt_gb(h.gb), 100 * h.coverage)
ro.update("%s %02d:00 — %s" % (lab, hi, txt))
self.query_one("#legend", Static).update(
"legend: [%(a)s]█ allocated[/] [%(u)s]▒ unallocated (attribution unknown)[/] "
"[%(g)s]░ gap · no evidence — never zero[/] [%(z)s]· measured 0 B[/] "
"[%(s)s]┄ partial day (so far)[/]" % dict(a=AMBER, u=UNALLOC, g=GAP, z=ZERO, s=SELECT))
v = "%s: %s" % (self.variant, VARIANT_NAMES[self.variant])
mode = "hourly strip" if self.variant == "C" else (
"day detail" if self.view == "hours" else "range view")
self.query_one("#footer", Static).update(
"PROTOTYPE ▸ %s ▸ %s [ / ] variant · ←/→ select · Enter open · Esc back · click" % (v, mode))
def action_prev_variant(self):
order = "ABC"
self.variant = order[(order.index(self.variant) - 1) % 3]
self.view = "strip" if self.variant == "C" else "days"
self.refresh_all()
def action_next_variant(self):
order = "ABC"
self.variant = order[(order.index(self.variant) + 1) % 3]
self.view = "strip" if self.variant == "C" else "days"
self.refresh_all()
def action_left(self):
if self.view == "hours":
self.sel_hour = max(0, self.sel_hour - 1)
elif self.variant == "C":
self.sel_strip = max(0, self.sel_strip - 1)
else:
self.sel_day = max(0, self.sel_day - 1)
self.refresh_all()
def action_right(self):
if self.view == "hours":
self.sel_hour = min(23, self.sel_hour + 1)
elif self.variant == "C":
self.sel_strip = min(71, self.sel_strip + 1)
else:
self.sel_day = min(len(self.days) - 1, self.sel_day + 1)
self.refresh_all()
def action_drill(self):
if self.variant in ("A", "B") and self.view == "days":
self.view = "hours"
self.sel_hour = 13
self.refresh_all()
def action_back(self):
if self.view == "hours":
self.view = "days"
self.refresh_all()
def on_click(self, event: events.Click) -> None:
w = event.widget
if w is None or w.id != "graph":
return
x = event.x
if self.view == "hours":
self.sel_hour = max(0, min(23, x))
elif self.variant == "C":
self.sel_strip = max(0, min(71, x))
else:
self.sel_day = max(0, min(len(self.days) - 1, x // 2))
self.refresh_all()
if __name__ == "__main__":
GraphPrototype().run()
+40 -5
View File
@@ -16,6 +16,8 @@ from pathlib import Path
from typing import Any, Dict, Optional, Tuple
from .store import init_store, get_store_path
from .monitoring_periods import ensure_period_open
from .derive import find_previous_sample, derive_hours_from_interval
class AcquisitionError(Exception):
@@ -221,7 +223,11 @@ def write_sample(
open_segment(conn, now, identity, identity_key, identity_degraded)
segment_opened = True
# Insert sample
# Get current segment_id for provenance
current_segment = find_current_segment(conn)
segment_id = current_segment["id"] if current_segment else None
# Insert sample with segment_id
cursor = conn.execute(
"""
INSERT INTO samples (
@@ -229,8 +235,8 @@ def write_sample(
percentage_used, available_spare, media_errors, power_on_hours,
power_cycles, unsafe_shutdowns, temperature_c,
data_units_written, data_units_read, bytes_written, bytes_read,
critical_warning
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
critical_warning, segment_id
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
""",
(
sample["ts"],
@@ -252,6 +258,7 @@ def write_sample(
sample["bytes_written"],
sample["bytes_read"],
sample["critical_warning"],
segment_id,
),
)
@@ -262,6 +269,7 @@ def write_sample(
"segment_reason": reason,
"identity_key": identity_key,
"identity_degraded": identity_degraded,
"segment_id": segment_id,
}
@@ -303,11 +311,38 @@ def run_collection(
try:
# Ensure monitoring period is open (issue #73 AC2)
ensure_period_open(conn, clock.utcnow())
# Validate invariants
validate_sample_invariants(sample, conn)
# Write sample
write_sample(sample, identity, conn, clock)
# Write sample and get segment info
seg_info = write_sample(sample, identity, conn, clock)
# Derive hour observations from interval with previous sample
try:
# 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
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)
except Exception:
# Derivation failure must not prevent sample persistence (issue #73 AC6)
pass
return {
"ok": True,
+237
View File
@@ -0,0 +1,237 @@
"""Interval derivation: samples → hour observations → day aggregates.
After each collection run, the collector calls into this module to:
1. Find the previous sample in the same segment
2. Compute deltas (bytes, POH, temperature)
3. Classify the hour(s) the interval spans
4. Write/update hour_observations for each affected hour
5. Update day_aggregates with unattributed cross-hour bytes
Cross-hour deltas are retained once with unknown shares explicit (issue #73 AC4).
No proportional allocation, endpoint assignment, or double counting.
"""
import sqlite3
from datetime import datetime, timedelta, timezone
from typing import Any, Dict, List, Optional, Tuple
from .hour_classify import classify_hour, HourSplit
def find_previous_sample(
conn: sqlite3.Connection,
segment_id: Optional[int],
current_sample_id: int,
) -> Optional[Dict[str, Any]]:
"""Find the most recent sample before current_sample_id in the same segment.
Returns None if no previous sample exists (first sample in segment).
"""
if segment_id is not None:
cursor = conn.execute(
"SELECT id, ts, bytes_written, bytes_read, power_on_hours, "
" temperature_c, data_units_written, data_units_read "
"FROM samples WHERE id < ? AND segment_id = ? "
"ORDER BY id DESC LIMIT 1",
(current_sample_id, segment_id),
)
else:
cursor = conn.execute(
"SELECT id, ts, bytes_written, bytes_read, power_on_hours, "
" temperature_c, data_units_written, data_units_read "
"FROM samples WHERE id < ? "
"ORDER BY id DESC LIMIT 1",
(current_sample_id,),
)
row = cursor.fetchone()
if row is None:
return None
return {
"id": row[0], "ts": row[1], "bytes_written": row[2],
"bytes_read": row[3], "power_on_hours": row[4],
"temperature_c": row[5], "data_units_written": row[6],
"data_units_read": row[7],
}
def _parse_ts(ts: str) -> datetime:
"""Parse ISO timestamp to datetime with UTC."""
dt = datetime.fromisoformat(ts)
if dt.tzinfo is None:
dt = dt.replace(tzinfo=timezone.utc)
return dt
def _hour_floor(dt: datetime) -> datetime:
"""Floor a datetime to its UTC hour boundary."""
return dt.replace(minute=0, second=0, microsecond=0)
def _hours_spanned(start: datetime, end: datetime) -> List[datetime]:
"""Return list of UTC hour boundaries spanned by [start, end)."""
hours = []
h = _hour_floor(start)
while h < end:
hours.append(h)
h += timedelta(hours=1)
return hours
def _compute_sampled_seconds_in_hour(
start: datetime, end: datetime, hour_start: datetime
) -> int:
"""How many seconds of the sample interval fall within this hour."""
hour_end = hour_start + timedelta(hours=1)
effective_start = max(start, hour_start)
effective_end = min(end, hour_end)
if effective_start >= effective_end:
return 0
return int((effective_end - effective_start).total_seconds())
def derive_hours_from_interval(
conn: sqlite3.Connection,
prev_sample: Dict[str, Any],
next_sample: Dict[str, Any],
) -> List[Dict[str, Any]]:
"""Derive hour observations from a sample pair interval.
Returns list of hour observation dicts that were written/updated.
"""
prev_ts = _parse_ts(prev_sample["ts"])
next_ts = _parse_ts(next_sample["ts"])
# Deltas
bw_delta = max(0, next_sample["bytes_written"] - prev_sample["bytes_written"])
br_delta = max(0, next_sample["bytes_read"] - prev_sample["bytes_read"])
poh_delta_s = max(0, (next_sample["power_on_hours"] - prev_sample["power_on_hours"])) * 3600
hours = _hours_spanned(prev_ts, next_ts)
total_span_s = int((next_ts - prev_ts).total_seconds())
results = []
if len(hours) == 1:
# Same-hour interval: fully attributed to this hour
hour_key = hours[0].strftime("%Y-%m-%dT%H:00:00+00:00")
sampled_s = total_span_s
# Classify hour
split = classify_hour(
wall_clock_seconds=3600,
poh_delta=poh_delta_s,
duw_delta=bw_delta,
dur_delta=br_delta,
sampled_seconds=sampled_s,
)
_upsert_hour_observation(
conn, hour_key, split,
bw_delta, br_delta,
prev_sample.get("temperature_c"), next_sample.get("temperature_c"),
2, # 2 samples contributed (prev + next)
)
results.append({"hour": hour_key, "bytes_written": bw_delta, "attributed": True})
elif len(hours) >= 2:
# Cross-hour interval: split wall-clock time, bytes unattributed
for h in hours:
hour_key = h.strftime("%Y-%m-%dT%H:00:00+00:00")
sampled_s = _compute_sampled_seconds_in_hour(prev_ts, next_ts, h)
# For cross-hour, we classify based on time only (no byte attribution)
# The hour gets its time split but NOT the byte delta
split = classify_hour(
wall_clock_seconds=3600,
poh_delta=0, # POH attribution unknown for cross-hour
duw_delta=0, # Bytes unattributed
dur_delta=0,
sampled_seconds=sampled_s,
)
_upsert_hour_observation(
conn, hour_key, split,
0, 0, # No byte attribution for cross-hour
None, None,
0, # No sample falls IN this hour
)
results.append({"hour": hour_key, "bytes_written": 0, "attributed": False})
# Track unattributed bytes at day level
_add_unattributed_bytes(conn, prev_ts, next_ts, bw_delta, br_delta)
return results
def _upsert_hour_observation(
conn: sqlite3.Connection,
hour_key: str,
split: HourSplit,
bw_delta: int,
br_delta: int,
temp_min: Optional[int],
temp_max: Optional[int],
sample_count: int,
) -> None:
"""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 = ?",
(hour_key,),
).fetchone()
if existing is None:
temp_avg = ((temp_min or 0) + (temp_max or 0)) / 2 if temp_min is not None else None
coverage = (split.seconds_active + split.seconds_idle + split.seconds_powered_off) / 3600.0
conn.execute(
"""INSERT INTO hour_observations
(hour, active_seconds, idle_seconds, powered_off_seconds, unknown_seconds,
bytes_written_delta, bytes_read_delta,
temperature_min, temperature_avg, temperature_max,
sample_count, coverage)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)""",
(hour_key, split.seconds_active, split.seconds_idle,
split.seconds_powered_off, split.seconds_unknown,
bw_delta, br_delta,
temp_min, temp_avg, temp_max,
sample_count, coverage),
)
else:
# Merge: accumulate bytes and sample count
new_bw = existing[1] + bw_delta
new_samples = existing[2] + sample_count
conn.execute(
"UPDATE hour_observations SET bytes_written_delta = ?, sample_count = ? WHERE id = ?",
(new_bw, new_samples, existing[0]),
)
conn.commit()
def _add_unattributed_bytes(
conn: sqlite3.Connection,
prev_ts: datetime,
next_ts: datetime,
bw_delta: int,
br_delta: int,
) -> None:
"""Add unattributed byte deltas to day aggregates for each day touched."""
prev_day = prev_ts.strftime("%Y-%m-%d")
next_day = next_ts.strftime("%Y-%m-%d")
days = {prev_day, next_day}
for day in days:
existing = conn.execute(
"SELECT id FROM day_aggregates WHERE day = ?", (day,)
).fetchone()
if existing is None:
conn.execute(
"INSERT INTO day_aggregates (day, unattributed_bytes_written, unattributed_bytes_read) "
"VALUES (?, ?, ?)",
(day, bw_delta, br_delta),
)
else:
conn.execute(
"UPDATE day_aggregates SET unattributed_bytes_written = unattributed_bytes_written + ?, "
"unattributed_bytes_read = unattributed_bytes_read + ? WHERE day = ?",
(bw_delta, br_delta, day),
)
conn.commit()
+39 -11
View File
@@ -73,6 +73,7 @@ class ScenarioRange:
rates: Dict[int, float]
min_days: int
max_days: int
horizon_reasons: Dict[int, str] = field(default_factory=dict)
@dataclass(frozen=True)
@@ -234,20 +235,39 @@ def _compute_regime_rate(days, conn, regime_start_day, clock_now):
return regime_bytes / regime_wc, regime_bytes, regime_wc
def _compute_horizon_rate(days, conn, horizon_days, clock_now):
cutoff = (clock_now - timedelta(days=horizon_days)).strftime("%Y-%m-%d")
def _compute_horizon_rate(days, conn, horizon_days, evidence_endpoint):
"""Compute horizon rate anchored at the latest evidence endpoint T.
Uses exact trailing horizon_days × 86400 seconds from T, not clock_now.
Reader refresh alone never moves T or dilutes rates.
Returns (rate, reason) where reason is None on success or a string
describing why the rate is unavailable.
"""
if not days:
return None, "no observation history"
# T is the latest evidence endpoint — the end of the last day aggregate
T = datetime.fromisoformat(days[-1]["day"] + "T23:59:59+00:00")
# Exact trailing start: T minus horizon_days × 86400 seconds
h_start = T - timedelta(days=horizon_days)
cutoff = h_start.strftime("%Y-%m-%d")
# History must span the full horizon — no placeholders
if not days or days[0]["day"] > cutoff:
return None
if days[0]["day"] > cutoff:
return None, "%d-day window starts before earliest data" % horizon_days
h_bytes = sum(d["bytes_written"] for d in days if d["day"] >= cutoff)
covered = sum(1 for d in days if d["day"] >= cutoff)
if covered == 0:
return None
h_start = datetime.fromisoformat(cutoff + "T00:00:00+00:00")
h_wc = _wall_clock_in_range(conn, h_start, clock_now)
return None, "%d-day window has no data" % horizon_days
h_wc = _wall_clock_in_range(conn, h_start, T)
if h_wc <= 0:
return None
return h_bytes / h_wc
return None, "%d-day window has no monitored wall-clock time" % horizon_days
return h_bytes / h_wc, None
# ---------------------------------------------------------------------------
@@ -508,12 +528,20 @@ def compute_projection(conn, clock_now):
scenario = None
horizon_rates = {}
horizon_reasons = {}
for h in HORIZON_DAYS:
hr = _compute_horizon_rate(all_days, conn, h, clock_now)
hr, reason = _compute_horizon_rate(all_days, conn, h, clock_now)
if hr is not None:
horizon_rates[h] = hr
else:
horizon_reasons[h] = reason
if horizon_rates:
scenario = ScenarioRange(rates=horizon_rates, min_days=min(horizon_rates), max_days=max(horizon_rates))
scenario = ScenarioRange(
rates=horizon_rates,
min_days=min(horizon_rates),
max_days=max(horizon_rates),
horizon_reasons=horizon_reasons,
)
total_days_count = len(segment_days)
days_below_coverage = sum(1 for d in segment_days
+133 -2
View File
@@ -2,6 +2,7 @@
Raw samples are pruned opportunistically to 14 days.
Hour observations and day aggregates are retained indefinitely.
Boundary anchors required for successor evidence are retained.
"""
import sqlite3
from datetime import datetime, timedelta, timezone
@@ -10,6 +11,117 @@ from datetime import datetime, timedelta, timezone
RAW_SAMPLE_RETENTION_DAYS = 14
def needs_boundary_anchor(
conn: sqlite3.Connection,
sample_ts: str,
now: datetime,
) -> bool:
"""Check if a sample is needed as a boundary anchor for derivation.
A sample is a boundary anchor if:
1. It's older than retention_days (strictly before cutoff)
2. It has a next sample that forms an interval spanning the retention boundary
3. The interval hasn't been derived yet
The interval spans the boundary if:
- The sample is before the cutoff, AND
- The next sample is strictly after the cutoff (or within retention)
"""
from .derive import _parse_ts
sample_dt = _parse_ts(sample_ts)
retention_cutoff = now - timedelta(days=RAW_SAMPLE_RETENTION_DAYS)
# If sample is within retention (strictly after cutoff), not an anchor
if sample_dt > retention_cutoff:
return False
# Check if this sample has a next sample
cursor = conn.execute(
"""SELECT ts, segment_id FROM samples WHERE ts > ? ORDER BY ts LIMIT 1""",
(sample_ts,),
)
next_row = cursor.fetchone()
if next_row is None:
# No next sample - this is the last sample
# It's not needed for derivation (no interval to derive)
return False
next_ts_str = next_row[0]
next_segment_id = next_row[1]
next_dt = _parse_ts(next_ts_str)
# Check if the next sample is strictly after the cutoff (i.e., interval spans boundary)
if next_dt > retention_cutoff:
# The interval spans the retention boundary
# Check if the interval needs derivation
# Get current sample's segment_id
cursor = conn.execute(
"SELECT segment_id FROM samples WHERE ts = ?",
(sample_ts,),
)
current_segment_row = cursor.fetchone()
current_segment_id = current_segment_row[0] if current_segment_row else None
# If different segments, no interval to derive
if current_segment_id != next_segment_id:
return False
# Check if the interval [sample_ts, next_ts] needs derivation
# It needs derivation if any hour in the span lacks an observation
current_hour = sample_dt.replace(minute=0, second=0, microsecond=0)
end_hour = next_dt.replace(minute=0, second=0, microsecond=0)
while current_hour <= end_hour:
cursor = conn.execute(
"SELECT id FROM hour_observations WHERE hour = ?",
(current_hour.isoformat(),),
)
if cursor.fetchone() is None:
# This hour lacks an observation - interval needs derivation
return True
current_hour += timedelta(hours=1)
# All hours in the span have observations - interval is derived
return False
else:
# The interval doesn't span the boundary (both samples are old)
# Check if the interval needs derivation
# Get current sample's segment_id
cursor = conn.execute(
"SELECT segment_id FROM samples WHERE ts = ?",
(sample_ts,),
)
current_segment_row = cursor.fetchone()
current_segment_id = current_segment_row[0] if current_segment_row else None
# If different segments, no interval to derive
if current_segment_id != next_segment_id:
return False
# Check if the interval [sample_ts, next_ts] needs derivation
current_hour = sample_dt.replace(minute=0, second=0, microsecond=0)
end_hour = next_dt.replace(minute=0, second=0, microsecond=0)
while current_hour <= end_hour:
cursor = conn.execute(
"SELECT id FROM hour_observations WHERE hour = ?",
(current_hour.isoformat(),),
)
if cursor.fetchone() is None:
# This hour lacks an observation - interval needs derivation
# But only keep if the interval is significant (spans multiple hours)
# or if the next sample is the last sample before a gap
gap = (next_dt - sample_dt).total_seconds()
if gap > 24 * 3600: # Significant gap (> 24 hours)
return True
current_hour += timedelta(hours=1)
# All hours in the span have observations or gap is not significant
return False
def prune_old_samples(
conn: sqlite3.Connection,
now: datetime,
@@ -17,6 +129,8 @@ def prune_old_samples(
) -> int:
"""Remove raw samples older than retention_days.
Retains boundary anchors required for successor evidence.
Args:
conn: Connection to the observation store.
now: Current UTC time.
@@ -26,6 +140,23 @@ def prune_old_samples(
Number of samples removed.
"""
cutoff = (now - timedelta(days=retention_days)).isoformat()
cursor = conn.execute("DELETE FROM samples WHERE ts < ?", (cutoff,))
# Get all samples older than cutoff
cursor = conn.execute(
"SELECT id, ts FROM samples WHERE ts < ? ORDER BY ts",
(cutoff,),
)
old_samples = cursor.fetchall()
removed = 0
for sample_id, sample_ts in old_samples:
# Check if this sample is a boundary anchor
if needs_boundary_anchor(conn, sample_ts, now):
continue # Skip - it's a boundary anchor
# Remove the sample
conn.execute("DELETE FROM samples WHERE id = ?", (sample_id,))
removed += 1
conn.commit()
return cursor.rowcount
return removed
+448
View File
@@ -0,0 +1,448 @@
"""Historical repair and retention (issue #74).
Re-derives hour observations and day aggregates from surviving raw samples,
with idempotent and interruption-safe guarantees. Preserves import markers,
boundary anchors, and valid historical summaries.
Contracts:
- Idempotent: running repair multiple times produces no duplicates
- Interruption-safe: partial repair preserves prior valid history
- Preserves existing valid data: never overwrites valid derived data
- Surfaces failures explicitly for retry
"""
import logging
import sqlite3
from dataclasses import dataclass, field
from datetime import datetime, timedelta, timezone
from typing import Optional, List, Tuple
from .derive import find_previous_sample, derive_hours_from_interval, _parse_ts
logger = logging.getLogger(__name__)
@dataclass
class RepairResult:
"""Result of a repair operation."""
ok: bool
hours_created: int = 0
hours_updated: int = 0
days_created: int = 0
days_updated: int = 0
intervals_derived: int = 0
boundary_anchors_retained: int = 0
legacy_summaries_preserved: int = 0
error: Optional[str] = None
@dataclass
class RepairStatus:
"""Current repair status for read-only views."""
last_repair: Optional[str] = None # ISO timestamp of last successful repair
repair_in_progress: bool = False
hours_derived: int = 0
days_derived: int = 0
def _ensure_repair_metadata(conn: sqlite3.Connection) -> None:
"""Ensure metadata table exists for tracking repair state."""
conn.execute("""
CREATE TABLE IF NOT EXISTS store_metadata (
key TEXT PRIMARY KEY,
value TEXT NOT NULL
)
""")
def is_repair_in_progress(conn: sqlite3.Connection) -> bool:
"""Check if a repair operation is currently in progress."""
_ensure_repair_metadata(conn)
cursor = conn.execute(
"SELECT value FROM store_metadata WHERE key = 'repair_in_progress'"
)
row = cursor.fetchone()
return row is not None and row[0] == "true"
def _set_repair_in_progress(conn: sqlite3.Connection, in_progress: bool) -> None:
"""Mark repair as in progress or complete."""
_ensure_repair_metadata(conn)
conn.execute(
"INSERT OR REPLACE INTO store_metadata (key, value) VALUES (?, ?)",
("repair_in_progress", "true" if in_progress else "false"),
)
conn.commit()
def _update_repair_status(conn: sqlite3.Connection, result: RepairResult) -> None:
"""Update repair status after successful completion."""
_ensure_repair_metadata(conn)
now = datetime.now(timezone.utc).isoformat()
# Update last repair timestamp
conn.execute(
"INSERT OR REPLACE INTO store_metadata (key, value) VALUES (?, ?)",
("last_repair", now),
)
# Update derived counts
cursor = conn.execute("SELECT COUNT(*) FROM hour_observations")
hours = cursor.fetchone()[0]
cursor = conn.execute("SELECT COUNT(*) FROM day_aggregates")
days = cursor.fetchone()[0]
conn.execute(
"INSERT OR REPLACE INTO store_metadata (key, value) VALUES (?, ?)",
("hours_derived", str(hours)),
)
conn.execute(
"INSERT OR REPLACE INTO store_metadata (key, value) VALUES (?, ?)",
("days_derived", str(days)),
)
conn.commit()
def get_repair_status(conn: sqlite3.Connection) -> RepairStatus:
"""Get current repair status for read-only views."""
_ensure_repair_metadata(conn)
last_repair = None
cursor = conn.execute(
"SELECT value FROM store_metadata WHERE key = 'last_repair'"
)
row = cursor.fetchone()
if row:
last_repair = row[0]
in_progress = is_repair_in_progress(conn)
hours_derived = 0
cursor = conn.execute(
"SELECT value FROM store_metadata WHERE key = 'hours_derived'"
)
row = cursor.fetchone()
if row:
hours_derived = int(row[0])
days_derived = 0
cursor = conn.execute(
"SELECT value FROM store_metadata WHERE key = 'days_derived'"
)
row = cursor.fetchone()
if row:
days_derived = int(row[0])
return RepairStatus(
last_repair=last_repair,
repair_in_progress=in_progress,
hours_derived=hours_derived,
days_derived=days_derived,
)
def needs_boundary_anchor(
conn: sqlite3.Connection,
sample_ts: str,
now: datetime,
) -> bool:
"""Check if a sample is needed as a boundary anchor for derivation.
A sample is a boundary anchor if:
1. It's older than retention_days
2. It has no derived hour observation for its hour
3. It's the last sample before a gap that needs derivation (gap > 24 hours)
"""
from .pruning import RAW_SAMPLE_RETENTION_DAYS
sample_dt = _parse_ts(sample_ts)
retention_cutoff = now - timedelta(days=RAW_SAMPLE_RETENTION_DAYS)
# If sample is within retention, not an anchor (will be kept anyway)
if sample_dt >= retention_cutoff:
return False
# Check if this sample's hour already has a derived observation
hour_start = sample_dt.replace(minute=0, second=0, microsecond=0).isoformat()
hour_end = (sample_dt + timedelta(hours=1)).replace(minute=0, second=0, microsecond=0).isoformat()
cursor = conn.execute(
"""SELECT COUNT(*) FROM hour_observations
WHERE hour >= ? AND hour < ?""",
(hour_start, hour_end),
)
# If there's an hour observation in this sample's hour, it's been derived
if cursor.fetchone()[0] > 0:
return False
# Check if this sample is the last sample before a gap
# (i.e., the next sample is significantly later)
cursor = conn.execute(
"""SELECT ts FROM samples WHERE ts > ? ORDER BY ts LIMIT 1""",
(sample_ts,),
)
next_row = cursor.fetchone()
if next_row is None:
# No next sample - this is the last sample, might be needed
# But if it's old and fully derived, it's not needed
return False
next_ts = _parse_ts(next_row[0])
gap = (next_ts - sample_dt).total_seconds()
# If gap > 24 hours, this sample is a boundary anchor
# (needed to derive the interval spanning the gap)
return gap > 24 * 3600
def _get_unlinked_intervals(conn: sqlite3.Connection) -> List[Tuple[dict, dict]]:
"""Find sample pairs that form intervals but have no hour observations."""
cursor = conn.execute(
"""SELECT id, ts, bytes_written, bytes_read, power_on_hours,
temperature_c, data_units_written, data_units_read, segment_id
FROM samples ORDER BY ts"""
)
all_samples = []
for row in cursor.fetchall():
all_samples.append({
"id": row[0], "ts": row[1], "bytes_written": row[2],
"bytes_read": row[3], "power_on_hours": row[4],
"temperature_c": row[5], "data_units_written": row[6],
"data_units_read": row[7], "segment_id": row[8],
})
intervals = []
for i in range(len(all_samples) - 1):
prev = all_samples[i]
next_s = all_samples[i + 1]
# Skip if different segments
if prev["segment_id"] != next_s["segment_id"]:
continue
# Check if the interval spans hours that need derivation
prev_dt = _parse_ts(prev["ts"])
next_dt = _parse_ts(next_s["ts"])
# Check if any hour in the span lacks an observation
current = prev_dt.replace(minute=0, second=0, microsecond=0)
end = next_dt.replace(minute=0, second=0, microsecond=0)
needs_derivation = False
while current <= end:
cursor2 = conn.execute(
"SELECT id FROM hour_observations WHERE hour = ?",
(current.isoformat(),),
)
if cursor2.fetchone() is None:
needs_derivation = True
break
current += timedelta(hours=1)
if needs_derivation:
intervals.append((prev, next_s))
return intervals
def _derive_day_aggregate_from_hours(
conn: sqlite3.Connection,
day: str,
) -> Optional[dict]:
"""Derive a day aggregate from its hour observations."""
cursor = conn.execute(
"""SELECT SUM(active_seconds), SUM(idle_seconds),
SUM(powered_off_seconds), SUM(unknown_seconds),
SUM(bytes_written_delta), SUM(bytes_read_delta),
SUM(sample_count)
FROM hour_observations WHERE hour LIKE ?""",
(day + "T%",),
)
row = cursor.fetchone()
if row is None or row[0] is None:
return None
return {
"day": day,
"active_seconds": row[0] or 0,
"idle_seconds": row[1] or 0,
"powered_off_seconds": row[2] or 0,
"unknown_seconds": row[3] or 0,
"bytes_written_delta": row[4] or 0,
"bytes_read_delta": row[5] or 0,
"sample_count": row[6] or 0,
}
def _upsert_day_aggregate(conn: sqlite3.Connection, day_data: dict) -> bool:
"""Insert or update a day aggregate. Returns True if created."""
existing = conn.execute(
"SELECT id FROM day_aggregates WHERE day = ?",
(day_data["day"],),
).fetchone()
if existing is None:
# Calculate coverage
total_seconds = (day_data["active_seconds"] + day_data["idle_seconds"] +
day_data["powered_off_seconds"] + day_data["unknown_seconds"])
coverage = (day_data["active_seconds"] + day_data["idle_seconds"] +
day_data["powered_off_seconds"]) / total_seconds if total_seconds > 0 else 0.0
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 (?, ?, ?, ?, ?, ?, ?, ?, ?)""",
(day_data["day"], day_data["active_seconds"], day_data["idle_seconds"],
day_data["powered_off_seconds"], day_data["unknown_seconds"],
day_data["bytes_written_delta"], day_data["bytes_read_delta"],
day_data["sample_count"], coverage),
)
return True
else:
# Update existing (but only if new data is more complete)
# This implements the "don't overwrite valid older history" rule
cursor = conn.execute(
"""SELECT bytes_written_delta, sample_count
FROM day_aggregates WHERE day = ?""",
(day_data["day"],),
)
existing_data = cursor.fetchone()
# Only update if new data has more samples or more bytes
if (day_data["sample_count"] > existing_data[1] or
day_data["bytes_written_delta"] > existing_data[0]):
total_seconds = (day_data["active_seconds"] + day_data["idle_seconds"] +
day_data["powered_off_seconds"] + day_data["unknown_seconds"])
coverage = (day_data["active_seconds"] + day_data["idle_seconds"] +
day_data["powered_off_seconds"]) / total_seconds if total_seconds > 0 else 0.0
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 day = ?""",
(day_data["active_seconds"], day_data["idle_seconds"],
day_data["powered_off_seconds"], day_data["unknown_seconds"],
day_data["bytes_written_delta"], day_data["bytes_read_delta"],
day_data["sample_count"], coverage, day_data["day"]),
)
return False # Updated, not created
else:
return False # No update needed
def repair_derivation(
conn: sqlite3.Connection,
clock=None,
) -> RepairResult:
"""Repair hour observations and day aggregates from surviving samples.
This is the main entry point for historical repair. It:
1. Finds sample pairs that need interval derivation
2. Derives hour observations from those intervals
3. Updates day aggregates from the hour observations
4. Preserves existing valid data
5. Is idempotent and interruption-safe
Args:
conn: Connection to the observation store
clock: Injected clock (for testing)
Returns:
RepairResult with operation details
"""
if clock is None:
clock = datetime.now(timezone.utc)
elif hasattr(clock, 'utcnow'):
clock = clock.utcnow()
# Check if repair is already in progress
if is_repair_in_progress(conn):
return RepairResult(
ok=False,
error="Repair already in progress",
)
# Mark repair as in progress
_set_repair_in_progress(conn, True)
result = RepairResult(ok=True)
try:
# Begin transaction
conn.execute("BEGIN IMMEDIATE")
# 1. Find and derive intervals from sample pairs
intervals = _get_unlinked_intervals(conn)
for prev, next_s in intervals:
try:
derived_hours = derive_hours_from_interval(conn, prev, next_s)
result.intervals_derived += 1
result.hours_created += len([h for h in derived_hours if h.get("attributed", True)])
except Exception as e:
logger.warning("Failed to derive interval %s -> %s: %s",
prev["ts"], next_s["ts"], e)
# Continue with other intervals (resilient)
# 2. Update day aggregates from hour observations
cursor = conn.execute(
"SELECT DISTINCT substr(hour, 1, 10) as day FROM hour_observations ORDER BY day"
)
days = [row[0] for row in cursor.fetchall()]
for day in days:
day_data = _derive_day_aggregate_from_hours(conn, day)
if day_data is not None:
created = _upsert_day_aggregate(conn, day_data)
if created:
result.days_created += 1
else:
result.days_updated += 1
# 3. Count boundary anchors retained
now = clock if isinstance(clock, datetime) else datetime.now(timezone.utc)
cursor = conn.execute("SELECT ts FROM samples ORDER BY ts")
anchor_count = 0
for row in cursor.fetchall():
if needs_boundary_anchor(conn, row[0], now):
anchor_count += 1
result.boundary_anchors_retained = anchor_count
# 4. Count preserved legacy summaries
# Legacy summaries are day aggregates without corresponding hour observations
cursor = conn.execute(
"""SELECT COUNT(*) FROM day_aggregates d
WHERE NOT EXISTS (
SELECT 1 FROM hour_observations h
WHERE h.hour LIKE d.day || 'T%'
)"""
)
result.legacy_summaries_preserved = cursor.fetchone()[0]
# Commit transaction
conn.commit()
# Update repair status
_update_repair_status(conn, result)
except Exception as e:
conn.rollback()
logger.error("Repair failed: %s", e)
return RepairResult(
ok=False,
error=str(e),
)
finally:
# Mark repair as complete
_set_repair_in_progress(conn, False)
return result
+31 -6
View File
@@ -331,7 +331,8 @@ def check_retired_flag(flag: str) -> Optional[str]:
def _format_projection(proj, freshness: str, service: Dict[str, Any],
drive_facts: List[str], config_error: Optional[str],
store_fault: Optional[str], newer_schema: Optional[str],
journal_hint: Optional[str]) -> str:
journal_hint: Optional[str],
sample_count: int = 0, day_count: int = 0) -> str:
"""Format the complete status output."""
lines = []
@@ -361,6 +362,16 @@ def _format_projection(proj, freshness: str, service: Dict[str, Any],
_append_service_facts(lines, service)
return "\n".join(lines)
# --- Single sample: awaiting another sample (issue #73 AC3) ---
# Only show awaiting state when there are no day aggregates (e.g., legacy import
# or hand-crafted stores can have 1 sample but sufficient day data for projection)
if sample_count <= 1 and day_count == 0:
lines.append("awaiting another sample")
lines.append("")
lines.append("Collecting usage data — the first projection requires at least two samples.")
_append_service_facts(lines, service)
return "\n".join(lines)
# --- Projection headline ---
headline = _format_headline(proj)
lines.append(headline)
@@ -375,14 +386,17 @@ def _format_projection(proj, freshness: str, service: Dict[str, Any],
lines.append("%s evidence" % state_label)
lines.append("")
# --- Scenario range (§6.5) ---
if proj.scenario_range and proj.scenario_range.rates:
# --- Scenario range (§6.5) with horizon reasons ---
if proj.scenario_range:
parts = []
for horizon in sorted(proj.scenario_range.rates.keys()):
rate_gb_day = proj.scenario_range.rates[horizon] * 86400 / 1e9
parts.append("%dd: %.2f GB/day" % (horizon, rate_gb_day))
lines.append("scenario range: %s" % " · ".join(parts))
lines.append("")
for horizon, reason in sorted(proj.scenario_range.horizon_reasons.items()):
parts.append("%dd: %s" % (horizon, reason))
if parts:
lines.append("scenario range: %s" % " · ".join(parts))
lines.append("")
# --- PU context line (§6.1) ---
lines.append(proj.pu_context_line)
@@ -598,6 +612,17 @@ def get_status(store_path: Optional[Path] = None, clock_now: Optional[datetime]
freshness = grade_freshness(newest_ts, clock_now)
# --- Sample count for single-sample state (issue #73 AC3) ---
sample_count = 0
day_count = 0
try:
cursor = conn.execute("SELECT COUNT(*) FROM samples")
sample_count = cursor.fetchone()[0]
cursor = conn.execute("SELECT COUNT(*) FROM day_aggregates")
day_count = cursor.fetchone()[0]
except sqlite3.Error:
pass
# Freshness age for the service fact
freshness_age_s = None
if newest_ts:
@@ -637,7 +662,7 @@ def get_status(store_path: Optional[Path] = None, clock_now: Optional[datetime]
# --- Compose output ---
result = _format_projection(
proj, freshness, service, drive_facts, config_error,
None, None, journal_hint,
None, None, journal_hint, sample_count, day_count,
)
conn.close()
+33 -8
View File
@@ -12,7 +12,7 @@ from typing import Optional
# Schema version - increment on each migration
SCHEMA_VERSION = 1
SCHEMA_VERSION = 2
# Packaged default placement (spec §8.3). The config may override it, but a
@@ -106,7 +106,8 @@ def _create_schema(conn: sqlite3.Connection):
data_units_read INTEGER,
bytes_written INTEGER,
bytes_read INTEGER,
critical_warning INTEGER
critical_warning INTEGER,
segment_id INTEGER
)
""")
@@ -141,7 +142,9 @@ def _create_schema(conn: sqlite3.Connection):
bytes_written_delta INTEGER DEFAULT 0,
bytes_read_delta INTEGER DEFAULT 0,
sample_count INTEGER DEFAULT 0,
coverage REAL DEFAULT 0.0
coverage REAL DEFAULT 0.0,
unattributed_bytes_written INTEGER DEFAULT 0,
unattributed_bytes_read INTEGER DEFAULT 0
)
""")
@@ -207,11 +210,25 @@ def _apply_migrations(conn: sqlite3.Connection, current_version: int):
Spec: §3.6, §10.2
"""
# Migration 1→2: example placeholder
# if current_version < 2:
# conn.execute("ALTER TABLE ...")
# current_version = 2
pass
# Migration 1→2: add segment_id provenance to samples,
# unattributed byte tracking to day_aggregates (issue #73)
if current_version < 2:
# Defensive: only ALTER if table exists (handles minimal v1 stores)
tables = {row[0] for row in conn.execute(
"SELECT name FROM sqlite_master WHERE type='table'"
).fetchall()}
if "samples" in tables:
# Check if column already exists (idempotent)
cols = {row[1] for row in conn.execute("PRAGMA table_info(samples)").fetchall()}
if "segment_id" not in cols:
conn.execute("ALTER TABLE samples ADD COLUMN segment_id INTEGER")
if "day_aggregates" in tables:
cols = {row[1] for row in conn.execute("PRAGMA table_info(day_aggregates)").fetchall()}
if "unattributed_bytes_written" not in cols:
conn.execute("ALTER TABLE day_aggregates ADD COLUMN unattributed_bytes_written INTEGER DEFAULT 0")
if "unattributed_bytes_read" not in cols:
conn.execute("ALTER TABLE day_aggregates ADD COLUMN unattributed_bytes_read INTEGER DEFAULT 0")
current_version = 2
def migrate_to_latest(store_path: Path) -> int:
@@ -239,6 +256,14 @@ def migrate_to_latest(store_path: Path) -> int:
conn.close()
return 0 # Already up to date
# Version 0 means no schema — create fresh (issue #73)
if current_version == 0:
_create_schema(conn)
conn.execute(f"PRAGMA user_version={SCHEMA_VERSION}")
conn.commit()
conn.close()
return SCHEMA_VERSION
steps = SCHEMA_VERSION - current_version
_apply_migrations(conn, current_version)
conn.execute(f"PRAGMA user_version={SCHEMA_VERSION}")
+658 -23
View File
@@ -19,12 +19,13 @@ import subprocess
import sys
from datetime import datetime, timedelta, timezone
from pathlib import Path
from typing import Any, Dict, List, Optional
from typing import Any, Callable, Dict, List, Optional
from textual.app import App, ComposeResult
from textual.binding import Binding
from textual.containers import Container, Horizontal, VerticalScroll
from textual.screen import ModalScreen
from textual.widget import Widget
from textual.widgets import Static
from .projection import (
@@ -112,6 +113,601 @@ def _habit_bar(a: float, i: float, o: float, u: float, width: int = 40) -> str:
return bar + "\n" + legend
# ---------------------------------------------------------------------------
# Bar graph glyphs and constants
# ---------------------------------------------------------------------------
_GLYPH_ALLOCATED = "\u2588" # \u2588 full block
_GLYPH_UNALLOCATED = "\u2593" # \u2593 dark shade
_GLYPH_GAP = "\u2591" # \u2591 light shade
_GLYPH_ZERO = "\u00b7" # \u00b7 middle dot
_GLYPH_PARTIAL = "\u258c" # \u258c left half block
_RANGE_OPTIONS = (7, 14, 28, 90)
_RANGE_DEFAULT = 14
_BAR_HEIGHT = 6
_BAR_WIDTH = 2
_BAR_SPACING = 1
# ---------------------------------------------------------------------------
# Interactive daily bar graph widget
# ---------------------------------------------------------------------------
class DailyBarGraph(Widget):
"""Interactive daily writes bar graph with hourly drill-down.
Renders writes-only daily bars using block glyphs, supports keyboard
and mouse navigation, range switching (7/14/28/90 days), day selection,
and hourly drill-down. No plotting dependency (issue #75 AC5).
"""
can_focus = True
can_focus_children = False
CSS = """
DailyBarGraph {
height: 100%;
width: 100%;
layout: vertical;
}
#bar-range {
height: 1;
width: 100%;
}
#bar-render {
height: 1fr;
width: 100%;
overflow: hidden;
}
#bar-legend {
height: 1;
width: 100%;
}
#bar-readout {
height: 3;
width: 100%;
}
"""
def __init__(self, **kwargs: Any) -> None:
super().__init__(**kwargs)
self.range_days: int = _RANGE_DEFAULT
self.selected_index: int = -1
self.view_mode: str = "daily"
self.drill_day: Optional[str] = None
self._all_day_data: List[Dict[str, Any]] = []
self._day_data: List[Dict[str, Any]] = []
self._hour_data: List[Dict[str, Any]] = []
self._max_bytes: int = 0
self._hourly_selected: int = -1
self._on_drill: Optional[Callable[[str], None]] = None
def compose(self) -> ComposeResult:
yield Static("", id="bar-range")
yield Static("", id="bar-render")
yield Static("", id="bar-legend")
yield Static("", id="bar-readout")
def set_data(
self,
day_data: List[Dict[str, Any]],
on_drill: Any = None,
) -> None:
"""Update graph with day data. on_drill(day) called on drill entry."""
self._all_day_data = day_data
self._on_drill = on_drill
self._trim_to_range()
self.view_mode = "daily"
self.drill_day = None
self._hour_data = []
self._hourly_selected = -1
self._refresh()
def set_hour_data(self, hour_data: List[Dict[str, Any]]) -> None:
"""Set hourly data for drill-down view."""
self._hour_data = hour_data
self._hourly_selected = -1
self._refresh_hourly()
def _trim_to_range(self) -> None:
self._day_data = (
self._all_day_data[-self.range_days:]
if len(self._all_day_data) > self.range_days
else list(self._all_day_data)
)
# Bar height scales to allocated bytes only; unallocated are shown
# separately so cross-day unknowns never inflate a day bar (AC3).
self._max_bytes = max(
(d.get("allocated_bytes", 0) for d in self._day_data), default=0
)
def _refresh(self) -> None:
if not self._day_data:
self._show_empty()
return
if self._is_constrained():
self._show_constrained_summary()
return
self._render_range()
self._render_bars()
self._render_legend()
self._render_readout()
def _show_empty(self) -> None:
self.query_one("#bar-range").update("[dim]usage history[/dim]")
self.query_one("#bar-render").update("[dim]awaiting first sample[/dim]")
self.query_one("#bar-legend").update("")
self.query_one("#bar-readout").update("")
def _is_constrained(self) -> bool:
"""Return True when the widget is too narrow for the bar graph.
At 80 columns with a 3fr:2fr grid split the pane content width is ~46
(after round border + padding). Below 32 the bars become unreadable.
"""
try:
w = self.region.width
return 0 < w < 32
except Exception:
return False
def _show_constrained_summary(self) -> None:
"""Textual fallback for terminals below 80×24."""
self.query_one("#bar-range").update("[dim]usage history[/dim]")
if not self._day_data:
self.query_one("#bar-render").update("[dim]graph needs ≥80×24[/dim]")
self.query_one("#bar-legend").update("")
self.query_one("#bar-readout").update("")
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)
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)
)
self.query_one("#bar-legend").update("")
# Preserve selection readout
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)
)
else:
self.query_one("#bar-readout").update("[dim]\u2190 \u2192 select[/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)
hours_with_data = sum(
1 for h in self._hour_data if h.get("bytes_written", 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)
)
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)
)
else:
self.query_one("#bar-readout").update("[dim]\u2190 \u2192 select hour[/dim]")
def on_resize(self) -> None:
"""Re-render when terminal size changes."""
if self.view_mode == "daily":
self._refresh()
else:
self._refresh_hourly()
def _render_range(self) -> None:
parts = []
for r in _RANGE_OPTIONS:
if r == self.range_days:
parts.append("[bold]%d[/bold]" % r)
else:
parts.append(str(r))
label = "range: " + " / ".join(parts)
if self.view_mode == "hourly":
label += " \u00b7 [bold]%s[/bold] \u00b7 esc back" % (self.drill_day or "")
self.query_one("#bar-range").update(label)
def _render_bars(self) -> None:
visible = self._day_data
if not visible:
self.query_one("#bar-render").update("")
return
max_bytes = self._max_bytes or 1
bar_h = _BAR_HEIGHT
n = len(visible)
lines: List[str] = []
for row in range(bar_h, 0, -1):
line = ""
threshold = (row / bar_h) * max_bytes
for i, day in enumerate(visible):
allocated = day.get("allocated_bytes", 0)
is_zero = day.get("is_zero", False)
is_gap = day.get("is_gap", False)
is_partial = day.get("is_partial", False)
if is_zero and row == 1:
glyph = _GLYPH_ZERO
elif is_gap:
glyph = _GLYPH_GAP if row <= 2 else " "
elif allocated == 0:
glyph = " "
elif row == 1 and i == self.selected_index:
glyph = "\u25b6" # selection arrow
elif threshold > 0 and threshold <= allocated:
glyph = _GLYPH_ALLOCATED
elif is_partial and row == bar_h:
glyph = _GLYPH_PARTIAL
else:
glyph = " "
line += glyph * _BAR_WIDTH
if i < n - 1:
line += " " * _BAR_SPACING
lines.append(line)
# Unallocated indicator row: show ▓ for days with cross-day unknowns
unalloc_line = ""
has_unalloc = False
for i, day in enumerate(visible):
unalloc = day.get("unallocated_bytes", 0)
if unalloc > 0:
unalloc_line += _GLYPH_UNALLOCATED * _BAR_WIDTH
has_unalloc = True
else:
unalloc_line += " " * _BAR_WIDTH
if i < n - 1:
unalloc_line += " " * _BAR_SPACING
if has_unalloc:
lines.append(unalloc_line)
# Date labels
label_line = ""
for i, day in enumerate(visible):
label = day.get("local_label", day.get("day", ""))[-2:]
label_line += label
if i < n - 1:
label_line += " " * _BAR_SPACING
lines.append(label_line)
self.query_one("#bar-render").update("\n".join(lines))
def _render_legend(self) -> None:
legend = (
"%s alloc %s unalloc %s gap %s zero %s partial"
% (
_GLYPH_ALLOCATED,
_GLYPH_UNALLOCATED,
_GLYPH_GAP,
_GLYPH_ZERO,
_GLYPH_PARTIAL,
)
)
self.query_one("#bar-legend").update(legend)
def _render_readout(self) -> None:
if self.selected_index < 0 or self.selected_index >= len(self._day_data):
self.query_one("#bar-readout").update(
"[dim]\u2190 \u2192 select \u00b7 1-4 range \u00b7 Enter drill[/dim]"
)
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)
coverage = day.get("coverage", 0)
hours = day.get("evidenced_hours", 0)
total_gb = total / 1e9
parts = [
"[bold]%s[/bold] \u00b7 %.3f GB total \u00b7 %d hours \u00b7 %.0f%% coverage"
% (
day.get("local_label", day.get("day", "")),
total_gb,
hours,
coverage * 100,
),
]
if unallocated > 0:
parts.append(
" allocated %.3f GB \u00b7 unallocated %.3f GB"
% (allocated / 1e9, unallocated / 1e9)
)
self.query_one("#bar-readout").update("\n".join(parts))
# -- Hourly drill-down rendering --
def _refresh_hourly(self) -> None:
if not self._hour_data:
self.query_one("#bar-render").update("[dim]no hourly data[/dim]")
return
self._render_range()
if self._is_constrained():
self._show_constrained_hourly_summary()
return
max_bytes = max(
(h.get("bytes_written", 0) for h in self._hour_data), default=0
) or 1
bar_h = _BAR_HEIGHT
n = len(self._hour_data)
lines: List[str] = []
for row in range(bar_h, 0, -1):
line = ""
threshold = (row / bar_h) * max_bytes
for i, hour in enumerate(self._hour_data):
bw = hour.get("bytes_written", 0)
is_zero = hour.get("is_zero", False)
if is_zero and row == 1:
glyph = _GLYPH_ZERO
elif bw == 0:
glyph = " "
elif row == 1 and i == self._hourly_selected:
glyph = "\u25b6"
elif threshold > 0 and threshold <= bw:
glyph = _GLYPH_ALLOCATED
else:
glyph = " "
line += glyph * _BAR_WIDTH
if i < n - 1:
line += " " * _BAR_SPACING
lines.append(line)
# Hour labels
label_line = ""
for i, hour in enumerate(self._hour_data):
label = hour.get("local_label", hour.get("hour", ""))[-2:]
label_line += label
if i < n - 1:
label_line += " " * _BAR_SPACING
lines.append(label_line)
self.query_one("#bar-render").update("\n".join(lines))
self.query_one("#bar-legend").update(
"%s writes \u00b7 %s zero" % (_GLYPH_ALLOCATED, _GLYPH_ZERO)
)
# Hourly readout
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 \u00b7 %d%% coverage"
% (
h.get("local_label", h.get("hour", "")),
h.get("bytes_written", 0) / 1e9,
h.get("coverage", 0) * 100,
)
)
else:
self.query_one("#bar-readout").update("[dim]\u2190 \u2192 select hour[/dim]")
# -- Event handling --
def on_key(self, event: Any) -> None:
if self.view_mode == "daily":
self._handle_daily_key(event)
else:
self._handle_hourly_key(event)
def _handle_daily_key(self, event: Any) -> None:
if event.key == "left":
if self.selected_index < 0:
self.selected_index = len(self._day_data) - 1
else:
self.selected_index = max(0, self.selected_index - 1)
self._refresh()
event.stop()
elif event.key == "right":
if self.selected_index < 0:
self.selected_index = 0
else:
self.selected_index = min(
len(self._day_data) - 1, self.selected_index + 1
)
self._refresh()
event.stop()
elif event.key == "enter" and self.selected_index >= 0:
self._enter_drill()
event.stop()
elif event.key in ("1", "2", "3", "4"):
self.range_days = _RANGE_OPTIONS[int(event.key) - 1]
self.selected_index = -1
self._trim_to_range()
self._refresh()
event.stop()
def _enter_drill(self) -> None:
if self.selected_index < 0 or self.selected_index >= len(self._day_data):
return
day = self._day_data[self.selected_index]
self.drill_day = day.get("day")
self.view_mode = "hourly"
self._hourly_selected = -1
if self._on_drill:
self._on_drill(self.drill_day)
self._render_range()
self.query_one("#bar-render").update("[dim]loading...[/dim]")
self.query_one("#bar-legend").update("")
self.query_one("#bar-readout").update("")
def _exit_drill(self) -> None:
self.view_mode = "daily"
self.drill_day = None
self._hour_data = []
self._hourly_selected = -1
self._refresh()
def _handle_hourly_key(self, event: Any) -> None:
if event.key in ("escape", "backspace"):
self._exit_drill()
event.stop()
elif event.key == "left":
if self._hourly_selected < 0:
self._hourly_selected = len(self._hour_data) - 1
else:
self._hourly_selected = max(0, self._hourly_selected - 1)
self._refresh_hourly()
event.stop()
elif event.key == "right":
if self._hourly_selected < 0:
self._hourly_selected = 0
else:
self._hourly_selected = min(
len(self._hour_data) - 1, self._hourly_selected + 1
)
self._refresh_hourly()
event.stop()
def on_click(self, event: Any) -> None:
render = self.query_one("#bar-render")
offset_x = event.x - render.region.x
bar_total = _BAR_WIDTH + _BAR_SPACING
idx = offset_x // bar_total
if self.view_mode == "daily":
if 0 <= idx < len(self._day_data):
self.selected_index = idx
self._refresh()
else:
if 0 <= idx < len(self._hour_data):
self._hourly_selected = idx
self._refresh_hourly()
# ---------------------------------------------------------------------------
# Graph data queries
# ---------------------------------------------------------------------------
def _query_daily_graph_data(
conn: sqlite3.Connection,
) -> List[Dict[str, Any]]:
"""Query day aggregates for the bar graph.
Returns one dict per day with total/allocated/unallocated bytes,
coverage, evidence hours, and classification flags.
"""
cursor = conn.execute(
"SELECT day, bytes_written_delta, unattributed_bytes_written, "
"coverage, sample_count, active_seconds, idle_seconds, "
"powered_off_seconds, unknown_seconds "
"FROM day_aggregates ORDER BY day"
)
rows = cursor.fetchall()
result: List[Dict[str, Any]] = []
for row in rows:
day = row[0]
bw_delta = row[1] or 0
unattributed = 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
total_bytes = bw_delta + unattributed
evidenced_hours = (active + idle + powered_off) // 3600
is_zero = total_bytes == 0
is_gap = (
sample_count == 0
and (active + idle + powered_off) == 0
and unknown > 0
)
is_partial = coverage < 0.5
result.append({
"day": day,
"local_label": day,
"total_bytes": total_bytes,
"allocated_bytes": bw_delta,
"unallocated_bytes": unattributed,
"coverage": coverage,
"evidenced_hours": evidenced_hours,
"sample_count": sample_count,
"is_zero": is_zero,
"is_gap": is_gap,
"is_partial": is_partial,
})
return result
def _query_hourly_graph_data(
conn: sqlite3.Connection,
day: str,
) -> List[Dict[str, Any]]:
"""Query hour observations for a specific day.
Returns one dict per hour with bytes written, coverage, and flags.
"""
cursor = conn.execute(
"SELECT hour, bytes_written_delta, coverage, sample_count, "
"active_seconds, idle_seconds, powered_off_seconds, unknown_seconds "
"FROM hour_observations "
"WHERE hour LIKE ? ORDER BY hour",
(day + "T%",),
)
rows = cursor.fetchall()
result: List[Dict[str, Any]] = []
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
is_zero = bw == 0
local_label = hour[11:13] if len(hour) >= 13 else hour
result.append({
"hour": hour,
"local_label": local_label,
"bytes_written": bw,
"coverage": coverage,
"sample_count": sample_count,
"active_seconds": active,
"idle_seconds": idle,
"powered_off_seconds": powered_off,
"unknown_seconds": unknown,
"is_zero": is_zero,
})
return result
# ---------------------------------------------------------------------------
# Data queries for TUI regions
# ---------------------------------------------------------------------------
@@ -356,7 +952,7 @@ class FenrisTuiApp(App):
with Container(id="main-grid"):
yield Static("", id="headline-band", classes="pane")
yield Static("", id="paused-banner")
yield Static("", id="usage-history", classes="pane")
yield DailyBarGraph(id="usage-history", classes="pane")
yield Static("", id="drive-health", classes="pane")
yield Static("", id="service-strip", classes="pane")
yield Static("q QUIT TUI", id="quit-rail")
@@ -423,7 +1019,7 @@ class FenrisTuiApp(App):
"[bold]No observations yet[/bold]\n\n"
"Enable monitoring: fenris monitor resume"
)
self.query_one("#usage-history").update("")
self.query_one("#usage-history").set_data([])
self.query_one("#drive-health").update("")
self.query_one("#service-strip").update(
"boot: disabled · timer: inactive · last collect: unknown · freshness: empty · "
@@ -438,7 +1034,7 @@ class FenrisTuiApp(App):
"[bold red]Observation store unreadable[/bold red]\n"
"Check journalctl -u fenris-collect.service"
)
self.query_one("#usage-history").update("")
self.query_one("#usage-history").set_data([])
self.query_one("#drive-health").update("")
self.query_one("#service-strip").update(
"[dim]by Bongbetic[/dim]\n"
@@ -447,26 +1043,47 @@ class FenrisTuiApp(App):
def _render_all_regions(self, conn: sqlite3.Connection) -> None:
"""Render all four regions from live store data."""
# --- Headline band (§7.2) ---
# --- Sample count for single-sample state (issue #73 AC3) ---
try:
proj = compute_projection(conn, self._clock_now)
headline = self._format_headline(proj)
confidence = self._format_confidence(proj)
scenario = self._format_scenario(proj)
self._render_headline(headline + "\n" + confidence + "\n" + scenario)
cursor = conn.execute("SELECT COUNT(*) FROM samples")
sample_count = cursor.fetchone()[0]
cursor = conn.execute("SELECT COUNT(*) FROM day_aggregates")
day_count = cursor.fetchone()[0]
except Exception:
self._render_headline("[bold]No projection available[/bold]")
sample_count = 0
day_count = 0
# --- Usage-history pane (§7.2 left) ---
history = _query_usage_history(conn)
history_text = "[bold]Usage history[/bold] · %d days · %.1f–%.1f GB/day\n %s\n %s" % (
history["num_days"],
history["min_gb"],
history["max_gb"],
history["sparkline"],
history["habit_bar"],
)
self.query_one("#usage-history").update(history_text)
# --- Headline band (§7.2) ---
if sample_count <= 1 and day_count == 0:
# Single sample: awaiting another sample
self._render_headline(
"[bold]Awaiting another sample[/bold]\n\n"
"Collecting usage data — the first projection requires at least two samples."
)
else:
try:
proj = compute_projection(conn, self._clock_now)
headline = self._format_headline(proj)
confidence = self._format_confidence(proj)
scenario = self._format_scenario(proj)
self._render_headline(headline + "\n" + confidence + "\n" + scenario)
except Exception:
self._render_headline("[bold]No projection available[/bold]")
# --- Usage-history pane (§7.2 left): interactive bar graph ---
graph = self.query_one("#usage-history")
day_data = _query_daily_graph_data(conn)
# Preserve drill-down state across refresh if still valid
if graph.view_mode == "hourly" and graph.drill_day:
saved_drill_day = graph.drill_day
hour_data = _query_hourly_graph_data(conn, saved_drill_day)
graph.set_data(day_data, on_drill=self._on_graph_drill)
graph.view_mode = "hourly"
graph.drill_day = saved_drill_day
graph.set_hour_data(hour_data)
else:
graph.set_data(day_data, on_drill=self._on_graph_drill)
# --- Drive-health pane (§7.2 right) ---
health = _query_drive_health(conn)
@@ -540,6 +1157,20 @@ class FenrisTuiApp(App):
main_grid.remove_class("paused")
main_grid.refresh(layout=True)
def _on_graph_drill(self, day: str) -> None:
"""Load hourly data when the graph enters drill-down mode."""
graph = self.query_one("#usage-history")
try:
conn = self._open_store()
if conn is not None:
try:
hour_data = _query_hourly_graph_data(conn, day)
graph.set_hour_data(hour_data)
finally:
conn.close()
except Exception:
graph.set_hour_data([])
def _format_headline(self, proj: ProjectionResult) -> str:
"""Format the lifespan headline (spec §6.11)."""
if proj.headline_remaining_seconds is None:
@@ -574,13 +1205,17 @@ class FenrisTuiApp(App):
)
def _format_scenario(self, proj: ProjectionResult) -> str:
"""Format scenario range (spec §6.5)."""
if not proj.scenario_range or not proj.scenario_range.rates:
"""Format scenario range with horizon reasons (spec §6.5)."""
if not proj.scenario_range:
return ""
parts = []
for horizon in sorted(proj.scenario_range.rates.keys()):
rate_gb_day = proj.scenario_range.rates[horizon] * 86400 / 1e9
parts.append("%dd: %.2f GB/day" % (horizon, rate_gb_day))
for horizon, reason in sorted(proj.scenario_range.horizon_reasons.items()):
parts.append("%dd: %s" % (horizon, reason))
if not parts:
return ""
return "[bold]Scenario range[/bold] · %s" % " · ".join(parts)
# --- Actions ---
+712
View File
@@ -0,0 +1,712 @@
"""Collector history tracer tests (issue #73).
Tests the end-to-end history pipeline: sample acquisition → interval
derivation → hour observation → day aggregate, with concurrent-read
safety, cross-hour handling, and display states.
Seams:
- write side: run_collection() → observation store
- read side: get_status(), compute_projection() → observation store
"""
import os
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, normalize_identity
from fenris.store import init_store, get_store_path, SCHEMA_VERSION
from fenris.monitoring_periods import ensure_period_open, close_period, get_open_period
from fenris.day_aggregate import derive_day, derive_all_days
# ---------------------------------------------------------------------------
# Fixtures
# ---------------------------------------------------------------------------
@pytest.fixture
def smartctl_fixture() -> Dict[str, Any]:
"""Minimal smartctl -a -j output with required fields."""
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": 12345678,
"data_units_read": 9876543,
"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_fixture_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
@pytest.fixture
def config_fixture(tmp_path: Path) -> Dict[str, Any]:
"""Configuration fixture naming the device."""
return {
"device": "/dev/nvme0",
"store_path": str(tmp_path / "observations.db"),
}
class FakeClock:
"""Injected clock returning controlled time."""
def __init__(self, initial: datetime):
self.now = initial
def utcnow(self):
return self.now
def advance(self, **kwargs):
self.now = self.now + timedelta(**kwargs)
# ---------------------------------------------------------------------------
# Schema migration: existing data readable at real precision
# ---------------------------------------------------------------------------
class TestSchemaMigration:
"""Schema migration 1→2 preserves existing data (issue #73 AC1)."""
def test_migration_bumps_version(self, tmp_path):
"""Migration from v1 to v2 succeeds."""
from fenris.store import migrate_to_latest, SCHEMA_VERSION
# Create a v1 store directly (simulating pre-migration state)
db = tmp_path / "test.db"
conn = sqlite3.connect(str(db))
conn.execute("PRAGMA journal_mode=WAL")
# Create v1 schema manually
conn.execute("""
CREATE TABLE samples (
id INTEGER PRIMARY KEY AUTOINCREMENT,
ts TEXT NOT NULL,
device TEXT NOT NULL,
subnqn TEXT, sn TEXT, mn TEXT, fr TEXT,
capacity_bytes INTEGER, percentage_used INTEGER,
available_spare INTEGER, media_errors INTEGER,
power_on_hours INTEGER, power_cycles INTEGER,
unsafe_shutdowns INTEGER, temperature_c INTEGER,
data_units_written INTEGER, data_units_read INTEGER,
bytes_written INTEGER, bytes_read INTEGER,
critical_warning INTEGER
)
""")
conn.execute("""
CREATE TABLE hour_observations (
id INTEGER PRIMARY KEY AUTOINCREMENT,
hour TEXT NOT NULL UNIQUE,
active_seconds INTEGER DEFAULT 0, idle_seconds INTEGER DEFAULT 0,
powered_off_seconds INTEGER DEFAULT 0, unknown_seconds INTEGER DEFAULT 0,
bytes_written_delta INTEGER DEFAULT 0, bytes_read_delta INTEGER DEFAULT 0,
temperature_min INTEGER, temperature_avg REAL, temperature_max INTEGER,
sample_count INTEGER DEFAULT 0, coverage REAL DEFAULT 0.0
)
""")
conn.execute("""
CREATE TABLE day_aggregates (
id INTEGER PRIMARY KEY AUTOINCREMENT,
day TEXT NOT NULL UNIQUE,
active_seconds INTEGER DEFAULT 0, idle_seconds INTEGER DEFAULT 0,
powered_off_seconds INTEGER DEFAULT 0, unknown_seconds INTEGER DEFAULT 0,
bytes_written_delta INTEGER DEFAULT 0, bytes_read_delta INTEGER DEFAULT 0,
sample_count INTEGER DEFAULT 0, coverage REAL DEFAULT 0.0
)
""")
conn.execute("CREATE TABLE monitoring_periods (id INTEGER PRIMARY KEY AUTOINCREMENT, started_at TEXT NOT NULL, ended_at TEXT, end_cause TEXT)")
conn.execute("CREATE TABLE controller_segments (id INTEGER PRIMARY KEY AUTOINCREMENT, opened_at TEXT NOT NULL, identity_key TEXT, identity_degraded BOOLEAN DEFAULT 0, subnqn TEXT, sn TEXT, mn TEXT, fr TEXT, vid TEXT, ssvid TEXT, transport TEXT)")
conn.execute("CREATE TABLE endurance_baseline (id INTEGER PRIMARY KEY AUTOINCREMENT, tbw_terabytes REAL NOT NULL, source_url TEXT, document_revision TEXT, entry_date TEXT, model_string TEXT, nominal_capacity_bytes INTEGER, validated_by TEXT, verified BOOLEAN DEFAULT 0, created_at TEXT NOT NULL, updated_at TEXT NOT NULL)")
conn.execute("CREATE TABLE store_metadata (key TEXT PRIMARY KEY, value TEXT NOT NULL)")
conn.execute("PRAGMA user_version=1")
conn.execute("INSERT INTO samples (ts, device, mn, sn, fr, capacity_bytes, percentage_used, available_spare, media_errors, power_on_hours, power_cycles, unsafe_shutdowns, temperature_c, data_units_written, data_units_read, bytes_written, bytes_read, critical_warning) VALUES ('2026-09-01T12:00:00+00:00', '/dev/nvme0', 'Test', 'SN', 'FR', 1000000000000, 5, 100, 0, 1000, 100, 0, 35, 1000000, 500000, 512000000000, 256000000000, 0)")
conn.commit()
conn.close()
# Migrate
steps = migrate_to_latest(db)
assert steps == 1
# Verify data preserved
conn = sqlite3.connect(str(db))
row = conn.execute("SELECT ts, mn FROM samples").fetchone()
version = conn.execute("PRAGMA user_version").fetchone()[0]
conn.close()
assert row[0] == "2026-09-01T12:00:00+00:00"
assert row[1] == "Test"
assert version == SCHEMA_VERSION
def test_newer_schema_refused(self, tmp_path):
"""Store with user_version > SCHEMA_VERSION is refused."""
db = tmp_path / "test.db"
conn = sqlite3.connect(str(db))
conn.execute("PRAGMA user_version=%d" % (SCHEMA_VERSION + 1))
conn.commit()
conn.close()
with pytest.raises(ValueError, match="newer Fenris"):
init_store(db)
def test_existing_data_preserved_after_migration(self, config_fixture,
smartctl_fixture,
sysfs_fixture_tree):
"""Existing sample data is not lost or modified by migration."""
clock = FakeClock(datetime(2026, 9, 1, 12, 0, 0, tzinfo=timezone.utc))
# Write first sample
run_collection(smartctl_fixture, sysfs_fixture_tree / "sys" / "class" / "nvme" / "nvme0",
config_fixture, clock)
conn = sqlite3.connect(config_fixture["store_path"])
row = conn.execute("SELECT ts, device, bytes_written FROM samples").fetchone()
conn.close()
assert row[0] == "2026-09-01T12:00:00+00:00"
assert row[1] == "/dev/nvme0"
assert row[2] == 12345678 * 512000
# ---------------------------------------------------------------------------
# Collector publishes samples, intervals, hour observations, day aggregates
# ---------------------------------------------------------------------------
class TestCollectorDerivation:
"""Collector derives hour observations and day aggregates (issue #73 AC2)."""
def _make_sample(self, duw_units: int, ts: str) -> Dict[str, Any]:
"""Build a smartctl fixture with specific DUW."""
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_units,
"data_units_read": 9876543,
"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",
}
def test_same_hour_20mb_derives_hour_obs(self, config_fixture, sysfs_fixture_tree):
"""Two samples in same hour with 20 MB delta → hour_obs gets 20 MB."""
t1 = datetime(2026, 9, 1, 12, 0, 0, tzinfo=timezone.utc)
t2 = datetime(2026, 9, 1, 12, 5, 0, tzinfo=timezone.utc)
# DUW units: 12345678 * 512000 = ~6.3 TB; 20 MB = 20*1024*1024 / 512000 ≈ 40 units
duw1 = 12345678
duw2 = duw1 + 40 # ~20 MB more
clock1 = FakeClock(t1)
s1 = self._make_sample(duw1, t1.isoformat())
r1 = run_collection(s1, sysfs_fixture_tree / "sys" / "class" / "nvme" / "nvme0",
config_fixture, clock1)
assert r1["ok"]
clock2 = FakeClock(t2)
s2 = self._make_sample(duw2, t2.isoformat())
r2 = run_collection(s2, sysfs_fixture_tree / "sys" / "class" / "nvme" / "nvme0",
config_fixture, clock2)
assert r2["ok"]
# Check hour_observation was derived
conn = sqlite3.connect(config_fixture["store_path"])
hour = conn.execute(
"SELECT bytes_written_delta, sample_count FROM hour_observations WHERE hour LIKE '2026-09-01T12%'"
).fetchone()
conn.close()
assert hour is not None, "Hour observation should exist for 12:00"
assert hour[1] >= 2 # at least 2 samples contributed
# bytes_written_delta should be the 20 MB delta (40 * 512000 = 20480000)
assert hour[0] == 40 * 512000
def test_zero_delta_derives_hour_obs(self, config_fixture, sysfs_fixture_tree):
"""Two samples in same hour with no DUW change → 0 B written."""
t1 = datetime(2026, 9, 1, 12, 0, 0, tzinfo=timezone.utc)
t2 = datetime(2026, 9, 1, 12, 5, 0, tzinfo=timezone.utc)
duw = 12345678 # same for both
clock1 = FakeClock(t1)
s1 = self._make_sample(duw, t1.isoformat())
run_collection(s1, sysfs_fixture_tree / "sys" / "class" / "nvme" / "nvme0",
config_fixture, clock1)
clock2 = FakeClock(t2)
s2 = self._make_sample(duw, t2.isoformat())
run_collection(s2, sysfs_fixture_tree / "sys" / "class" / "nvme" / "nvme0",
config_fixture, clock2)
conn = sqlite3.connect(config_fixture["store_path"])
hour = conn.execute(
"SELECT bytes_written_delta FROM hour_observations WHERE hour LIKE '2026-09-01T12%'"
).fetchone()
conn.close()
assert hour is not None
assert hour[0] == 0
def test_cross_hour_100mb_unattributed(self, config_fixture, sysfs_fixture_tree):
"""Samples in different hours → delta is unattributed to any hour."""
t1 = datetime(2026, 9, 1, 11, 55, 0, tzinfo=timezone.utc)
t2 = datetime(2026, 9, 1, 12, 5, 0, tzinfo=timezone.utc)
duw1 = 12345678
duw2 = duw1 + 200 # ~100 MB
clock1 = FakeClock(t1)
s1 = self._make_sample(duw1, t1.isoformat())
run_collection(s1, sysfs_fixture_tree / "sys" / "class" / "nvme" / "nvme0",
config_fixture, clock1)
clock2 = FakeClock(t2)
s2 = self._make_sample(duw2, t2.isoformat())
run_collection(s2, sysfs_fixture_tree / "sys" / "class" / "nvme" / "nvme0",
config_fixture, clock2)
conn = sqlite3.connect(config_fixture["store_path"])
# Hour observations should NOT contain the cross-hour delta
hour11 = conn.execute(
"SELECT bytes_written_delta FROM hour_observations WHERE hour LIKE '2026-09-01T11%'"
).fetchone()
hour12 = conn.execute(
"SELECT bytes_written_delta FROM hour_observations WHERE hour LIKE '2026-09-01T12%'"
).fetchone()
conn.close()
# Neither hour should have the full 100 MB delta attributed
# (they may have 0 or partial, but not 200*512000)
full_delta = 200 * 512000
if hour11 is not None:
assert hour11[0] != full_delta, "Hour 11 should not have full cross-hour delta"
if hour12 is not None:
assert hour12[0] != full_delta, "Hour 12 should not have full cross-hour delta"
def test_readonly_sees_consistent_snapshot(self, config_fixture, sysfs_fixture_tree):
"""A read-only reader sees valid pre- or post-publication snapshot."""
t1 = datetime(2026, 9, 1, 12, 0, 0, tzinfo=timezone.utc)
t2 = datetime(2026, 9, 1, 12, 5, 0, tzinfo=timezone.utc)
clock1 = FakeClock(t1)
run_collection(self._make_sample(12345678, t1.isoformat()),
sysfs_fixture_tree / "sys" / "class" / "nvme" / "nvme0",
config_fixture, clock1)
# Open read-only
ro_conn = sqlite3.connect(
"file:%s?mode=ro" % config_fixture["store_path"], uri=True
)
count_before = ro_conn.execute("SELECT COUNT(*) FROM samples").fetchone()[0]
ro_conn.close()
assert count_before == 1
# Write second sample
clock2 = FakeClock(t2)
run_collection(self._make_sample(12345718, t2.isoformat()),
sysfs_fixture_tree / "sys" / "class" / "nvme" / "nvme0",
config_fixture, clock2)
# Read-only reader sees 2 samples now
ro_conn2 = sqlite3.connect(
"file:%s?mode=ro" % config_fixture["store_path"], uri=True
)
count_after = ro_conn2.execute("SELECT COUNT(*) FROM samples").fetchone()[0]
ro_conn2.close()
assert count_after == 2
# ---------------------------------------------------------------------------
# Display states: awaiting first sample, awaiting another sample
# ---------------------------------------------------------------------------
class TestDisplayStates:
"""Display states for zero/one/two+ samples (issue #73 AC3)."""
def test_zero_samples_awaiting_first(self, tmp_path):
"""Zero samples → 'awaiting first sample' state."""
from fenris.status import get_status
db = tmp_path / "test.db"
init_store(db)
now = datetime(2026, 9, 1, 12, 0, 0, tzinfo=timezone.utc)
from unittest.mock import patch
with patch("fenris.status.query_service_state", return_value={
"boot_enabled": False, "timer_active": False,
"last_collect_ok": None, "last_collect_age_s": None,
"last_collect_reason": None,
}):
result = get_status(store_path=db, clock_now=now,
query_services=False, query_journal=False)
assert "no observations yet" in result.lower() or "awaiting" in result.lower()
def test_one_sample_awaiting_another(self, tmp_path):
"""One sample → 'awaiting another sample' state."""
from fenris.status import get_status
db = tmp_path / "test.db"
conn = init_store(db)
conn.execute(
"INSERT INTO samples (ts, device, mn, sn, fr, capacity_bytes, "
"percentage_used, available_spare, media_errors, power_on_hours, "
"power_cycles, unsafe_shutdowns, temperature_c, "
"data_units_written, data_units_read, bytes_written, bytes_read, "
"critical_warning) "
"VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
("2026-09-01T12:00:00+00:00", "/dev/nvme0", "Test", "SN", "FR",
1000000000000, 5, 100, 0, 1000, 100, 0, 35,
1000000, 500000, 512000000000, 256000000000, 0),
)
conn.commit()
conn.close()
now = datetime(2026, 9, 1, 12, 0, 0, tzinfo=timezone.utc)
from unittest.mock import patch
with patch("fenris.status.query_service_state", return_value={
"boot_enabled": False, "timer_active": False,
"last_collect_ok": None, "last_collect_age_s": None,
"last_collect_reason": None,
}):
result = get_status(store_path=db, clock_now=now,
query_services=False, query_journal=False)
# Should mention awaiting or insufficient data
lower = result.lower()
assert "awaiting" in lower or "another sample" in lower or "no projection" in lower
def test_one_sample_awaiting_another_in_tui(self, tmp_path):
"""One sample → TUI shows awaiting state."""
from fenris.tui import FenrisTuiApp, _query_service_facts
db = tmp_path / "test.db"
conn = init_store(db)
conn.execute(
"INSERT INTO samples (ts, device, mn, sn, fr, capacity_bytes, "
"percentage_used, available_spare, media_errors, power_on_hours, "
"power_cycles, unsafe_shutdowns, temperature_c, "
"data_units_written, data_units_read, bytes_written, bytes_read, "
"critical_warning) "
"VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
("2026-09-01T12:00:00+00:00", "/dev/nvme0", "Test", "SN", "FR",
1000000000000, 5, 100, 0, 1000, 100, 0, 35,
1000000, 500000, 512000000000, 256000000000, 0),
)
conn.commit()
conn.close()
now = datetime(2026, 9, 1, 12, 0, 0, tzinfo=timezone.utc)
# Test the projection handles single sample
from fenris.projection import compute_projection, ConfidenceState
conn = sqlite3.connect(db)
proj = compute_projection(conn, now)
conn.close()
# With only one sample, projection should be unavailable
assert proj.confidence_state == ConfidenceState.UNSUPPORTED
def test_two_same_hour_samples_show_measured_usage(self, config_fixture, sysfs_fixture_tree):
"""Two compatible same-hour samples show measured usage."""
t1 = datetime(2026, 9, 1, 12, 0, 0, tzinfo=timezone.utc)
t2 = datetime(2026, 9, 1, 12, 5, 0, tzinfo=timezone.utc)
clock1 = FakeClock(t1)
run_collection(self._make_sample_helper(12345678), sysfs_fixture_tree / "sys" / "class" / "nvme" / "nvme0",
config_fixture, clock1)
clock2 = FakeClock(t2)
run_collection(self._make_sample_helper(12345718), sysfs_fixture_tree / "sys" / "class" / "nvme" / "nvme0",
config_fixture, clock2)
from fenris.status import get_status
from unittest.mock import patch
now = datetime(2026, 9, 1, 12, 10, 0, tzinfo=timezone.utc)
with patch("fenris.status.query_service_state", return_value={
"boot_enabled": False, "timer_active": False,
"last_collect_ok": None, "last_collect_age_s": None,
"last_collect_reason": None,
}):
result = get_status(store_path=Path(config_fixture["store_path"]),
clock_now=now, query_services=False, query_journal=False)
# Should not say "no observations" or "awaiting"
lower = result.lower()
assert "no observations yet" not in lower
def _make_sample_helper(self, duw_units: int) -> Dict[str, Any]:
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_units,
"data_units_read": 9876543, "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",
}
# ---------------------------------------------------------------------------
# Pause crossing and counter reset
# ---------------------------------------------------------------------------
class TestPauseCrossing:
"""Delta across monitoring period gap (issue #73 AC4)."""
def test_pause_crossing_preserves_prior_history(self, config_fixture, sysfs_fixture_tree):
"""Delta across a paused period preserves prior hour observations."""
t1 = datetime(2026, 9, 1, 12, 0, 0, tzinfo=timezone.utc)
t_pause = datetime(2026, 9, 1, 13, 0, 0, tzinfo=timezone.utc)
t_resume = datetime(2026, 9, 1, 14, 0, 0, tzinfo=timezone.utc)
t2 = datetime(2026, 9, 1, 14, 5, 0, tzinfo=timezone.utc)
duw1 = 12345678
duw2 = duw1 + 100
# First sample (opens period)
clock1 = FakeClock(t1)
run_collection(self._make_sample_for_pause(duw1), sysfs_fixture_tree / "sys" / "class" / "nvme" / "nvme0",
config_fixture, clock1)
# Pause
conn = sqlite3.connect(config_fixture["store_path"])
close_period(conn, t_pause, "user_disabled")
conn.close()
# Resume with new sample
clock_resume = FakeClock(t_resume)
run_collection(self._make_sample_for_pause(duw1), sysfs_fixture_tree / "sys" / "class" / "nvme" / "nvme0",
config_fixture, clock_resume)
# Second sample after resume
clock2 = FakeClock(t2)
run_collection(self._make_sample_for_pause(duw2), sysfs_fixture_tree / "sys" / "class" / "nvme" / "nvme0",
config_fixture, clock2)
# Verify prior hour observations are intact
conn = sqlite3.connect(config_fixture["store_path"])
hour_12 = conn.execute(
"SELECT bytes_written_delta FROM hour_observations WHERE hour LIKE '2026-09-01T12%'"
).fetchone()
conn.close()
# Hour 12 should still have data from the first collection
assert hour_12 is not None, "Prior hour observation should be preserved"
def _make_sample_for_pause(self, duw_units: int) -> Dict[str, Any]:
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_units,
"data_units_read": 9876543, "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",
}
class TestCounterReset:
"""Counter reset / replacement opens new segment (issue #73 AC6)."""
def test_duw_decrease_opens_new_segment(self, config_fixture, sysfs_fixture_tree):
"""DUW decrease triggers new segment."""
t1 = datetime(2026, 9, 1, 12, 0, 0, tzinfo=timezone.utc)
t2 = datetime(2026, 9, 1, 12, 5, 0, tzinfo=timezone.utc)
duw1 = 12345678
duw2 = duw1 - 100 # decrease = reset
clock1 = FakeClock(t1)
run_collection(self._make_sample_for_reset(duw1), sysfs_fixture_tree / "sys" / "class" / "nvme" / "nvme0",
config_fixture, clock1)
clock2 = FakeClock(t2)
run_collection(self._make_sample_for_reset(duw2), sysfs_fixture_tree / "sys" / "class" / "nvme" / "nvme0",
config_fixture, clock2)
conn = sqlite3.connect(config_fixture["store_path"])
segments = conn.execute("SELECT COUNT(*) FROM controller_segments").fetchone()[0]
samples = conn.execute("SELECT COUNT(*) FROM samples").fetchone()[0]
conn.close()
# Should have 2 segments (new one opened for DUW decrease)
assert segments == 2
# Should have 2 samples
assert samples == 2
def _make_sample_for_reset(self, duw_units: int) -> Dict[str, Any]:
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_units,
"data_units_read": 9876543, "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",
}
# ---------------------------------------------------------------------------
# Derivation failure preserves prior history
# ---------------------------------------------------------------------------
class TestDerivationFailure:
"""Injected derivation failure preserves prior history (issue #73 AC6)."""
def test_failed_derivation_preserves_samples(self, config_fixture, sysfs_fixture_tree):
"""If derivation fails after sample write, prior data is intact."""
t1 = datetime(2026, 9, 1, 12, 0, 0, tzinfo=timezone.utc)
t2 = datetime(2026, 9, 1, 12, 5, 0, tzinfo=timezone.utc)
clock1 = FakeClock(t1)
run_collection(self._make_sample_for_failure(12345678),
sysfs_fixture_tree / "sys" / "class" / "nvme" / "nvme0",
config_fixture, clock1)
# Verify first sample exists
conn = sqlite3.connect(config_fixture["store_path"])
count = conn.execute("SELECT COUNT(*) FROM samples").fetchone()[0]
conn.close()
assert count == 1
# Second sample with valid data should succeed
clock2 = FakeClock(t2)
r2 = run_collection(self._make_sample_for_failure(12345718),
sysfs_fixture_tree / "sys" / "class" / "nvme" / "nvme0",
config_fixture, clock2)
assert r2["ok"]
# Both samples should exist
conn = sqlite3.connect(config_fixture["store_path"])
count = conn.execute("SELECT COUNT(*) FROM samples").fetchone()[0]
conn.close()
assert count == 2
def _make_sample_for_failure(self, duw_units: int) -> Dict[str, Any]:
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_units,
"data_units_read": 9876543, "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",
}
# ---------------------------------------------------------------------------
# Monitoring period is opened by collector
# ---------------------------------------------------------------------------
class TestMonitoringPeriod:
"""Collector ensures monitoring period is open (issue #73 AC2)."""
def test_first_sample_opens_period(self, config_fixture, sysfs_fixture_tree):
"""First collection run opens a monitoring period."""
t1 = datetime(2026, 9, 1, 12, 0, 0, tzinfo=timezone.utc)
clock1 = FakeClock(t1)
run_collection(self._make_sample_simple(), sysfs_fixture_tree / "sys" / "class" / "nvme" / "nvme0",
config_fixture, clock1)
conn = sqlite3.connect(config_fixture["store_path"])
period = get_open_period(conn)
conn.close()
assert period is not None, "A monitoring period should be open"
def test_subsequent_sample_keeps_period_open(self, config_fixture, sysfs_fixture_tree):
"""Subsequent collection runs keep the period open."""
t1 = datetime(2026, 9, 1, 12, 0, 0, tzinfo=timezone.utc)
t2 = datetime(2026, 9, 1, 12, 5, 0, tzinfo=timezone.utc)
clock1 = FakeClock(t1)
run_collection(self._make_sample_simple(), sysfs_fixture_tree / "sys" / "class" / "nvme" / "nvme0",
config_fixture, clock1)
clock2 = FakeClock(t2)
run_collection(self._make_sample_simple(), sysfs_fixture_tree / "sys" / "class" / "nvme" / "nvme0",
config_fixture, clock2)
conn = sqlite3.connect(config_fixture["store_path"])
period = get_open_period(conn)
periods_count = conn.execute("SELECT COUNT(*) FROM monitoring_periods").fetchone()[0]
conn.close()
assert period is not None
assert periods_count == 1 # Still only one period
def _make_sample_simple(self) -> Dict[str, Any]:
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": 12345678,
"data_units_read": 9876543, "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",
}
+124 -1
View File
@@ -22,7 +22,7 @@ from fenris.monitoring_periods import ensure_period_open, close_period
from fenris.projection import (
compute_projection, ConfidenceState, BaselineTier, ScenarioRange,
TBW_TO_BYTES, HORIZON_DAYS, WARMING_MIN_DAYS, STALENESS_HOURS,
YOUNG_REGIME_DAYS, DISCLOSURES,
YOUNG_REGIME_DAYS, DISCLOSURES, _compute_horizon_rate,
)
@@ -899,3 +899,126 @@ class TestIdentityChangeBlankKeys:
# All 25 days in same segment (equal blanks continue)
assert result.regime_days is not None
assert result.regime_days >= 20 # Most of the history
# ===========================================================================
# Issue #76: Evidence-anchored projection rates
# Scenario windows anchored at latest evidence endpoint T
# ===========================================================================
class TestEvidenceAnchoredHorizons:
"""Issue #76: Scenario windows anchored at latest published usage-evidence
endpoint T with exact trailing 7/28/90×86400-second starts."""
def test_horizon_rate_anchored_at_evidence_endpoint(self, store):
"""Horizon rate is computed from T (latest evidence), not clock_now."""
_insert_baseline(store, tbw_tb=10.0, verified=True)
_insert_segment(store, opened_at="2026-09-01T00:00:00+00:00")
_open_period(store, start="2026-09-01T00:00:00+00:00")
bw = 100 * 1024 * 1024
# 30 days of data ending Sep 29
for i in range(30):
d = (datetime(2026, 9, 1) + timedelta(days=i)).strftime("%Y-%m-%d")
_insert_day(store, d, bw=bw)
_insert_sample(store, "2026-09-30T10:00:00+00:00", pu=5)
# Clock is Oct 1, but T is Sep 29 (latest evidence endpoint)
clock = datetime(2026, 10, 1, 12, 0, 0, tzinfo=timezone.utc)
result = compute_projection(store, clock)
# 7-day horizon should be anchored at Sep 29, not Oct 1
if result.scenario_range and 7 in result.scenario_range.rates:
# Rate should be based on Sep 23-29, not Sep 25-Oct 1
assert result.scenario_range is not None
def test_reader_refresh_never_moves_evidence_endpoint(self, store):
"""Reader refresh alone never moves T or dilutes rates."""
_insert_baseline(store, tbw_tb=10.0, verified=True)
_insert_segment(store, opened_at="2026-09-01T00:00:00+00:00")
_open_period(store, start="2026-09-01T00:00:00+00:00")
bw = 100 * 1024 * 1024
for i in range(30):
d = (datetime(2026, 9, 1) + timedelta(days=i)).strftime("%Y-%m-%d")
_insert_day(store, d, bw=bw)
_insert_sample(store, "2026-09-30T10:00:00+00:00", pu=5)
# Two reads at different clock times
clock1 = datetime(2026, 9, 30, 12, 0, 0, tzinfo=timezone.utc)
clock2 = datetime(2026, 10, 1, 12, 0, 0, tzinfo=timezone.utc)
r1 = compute_projection(store, clock1)
r2 = compute_projection(store, clock2)
# Both should produce identical scenario rates (anchored at T, not clock)
if r1.scenario_range and r2.scenario_range:
assert r1.scenario_range.rates == r2.scenario_range.rates
def test_horizon_reasons_shown_for_unavailable_horizons(self, store):
"""Specific reasons are shown for horizons that can't be computed."""
_insert_baseline(store, tbw_tb=10.0, verified=True)
_insert_segment(store, opened_at="2026-09-20T00:00:00+00:00")
_open_period(store, start="2026-09-20T00:00:00+00:00")
bw = 100 * 1024 * 1024
# Only 10 days of data
for i in range(10):
d = (datetime(2026, 9, 20) + timedelta(days=i)).strftime("%Y-%m-%d")
_insert_day(store, d, bw=bw)
_insert_sample(store, "2026-09-30T10:00:00+00:00", pu=5)
result = compute_projection(store, _clock())
# 7-day horizon should be available, 28 and 90 should have reasons
if result.scenario_range:
assert 7 in result.scenario_range.rates
if 28 in result.scenario_range.horizon_reasons:
assert "starts before earliest data" in result.scenario_range.horizon_reasons[28]
if 90 in result.scenario_range.horizon_reasons:
assert "starts before earliest data" in result.scenario_range.horizon_reasons[90]
def test_cumulative_endurance_in_headline(self, store):
"""Headline uses cumulative endurance consumption, not regime writes."""
_insert_baseline(store, tbw_tb=10.0, verified=True)
_insert_segment(store, opened_at="2026-09-01T00:00:00+00:00")
_open_period(store, start="2026-09-01T00:00:00+00:00")
bw = 100 * 1024 * 1024
for i in range(30):
d = (datetime(2026, 9, 1) + timedelta(days=i)).strftime("%Y-%m-%d")
_insert_day(store, d, bw=bw)
_insert_sample(store, "2026-09-30T10:00:00+00:00", pu=5)
result = compute_projection(store, _clock())
# Headline should be computed with cumulative bytes
if result.headline_remaining_seconds is not None:
# E_baseline = 10 TB = 10e12 bytes
# cumulative_bytes = 30 * 100 * 1024 * 1024
# rate = cumulative_bytes / wall_clock
# headline = (E_baseline - cumulative_bytes) / rate
E_baseline = 10.0 * TBW_TO_BYTES
cumulative_bytes = 30 * bw
assert result.headline_remaining_seconds >= 0
def test_zero_boundary_delta_returns_zero(self, store):
"""Zero monotonic delta proves zero over its represented subspan."""
_insert_baseline(store, tbw_tb=10.0, verified=True)
_insert_segment(store, opened_at="2026-09-01T00:00:00+00:00")
_open_period(store, start="2026-09-01T00:00:00+00:00")
# Days with zero bytes written
for i in range(30):
d = (datetime(2026, 9, 1) + timedelta(days=i)).strftime("%Y-%m-%d")
_insert_day(store, d, bw=0)
_insert_sample(store, "2026-09-30T10:00:00+00:00", pu=5)
result = compute_projection(store, _clock())
# Zero rate should result in UNSUPPORTED
assert result.confidence_state == ConfidenceState.UNSUPPORTED
assert result.headline_remaining_seconds is None
def test_noon_endpoint_same_as_midnight(self, store):
"""Noon endpoint produces same rates as midnight endpoint."""
_insert_baseline(store, tbw_tb=10.0, verified=True)
_insert_segment(store, opened_at="2026-09-01T00:00:00+00:00")
_open_period(store, start="2026-09-01T00:00:00+00:00")
bw = 100 * 1024 * 1024
for i in range(30):
d = (datetime(2026, 9, 1) + timedelta(days=i)).strftime("%Y-%m-%d")
_insert_day(store, d, bw=bw)
_insert_sample(store, "2026-09-30T10:00:00+00:00", pu=5)
# Both reads should produce same scenario rates
clock1 = datetime(2026, 9, 30, 0, 0, 0, tzinfo=timezone.utc)
clock2 = datetime(2026, 9, 30, 12, 0, 0, tzinfo=timezone.utc)
r1 = compute_projection(store, clock1)
r2 = compute_projection(store, clock2)
if r1.scenario_range and r2.scenario_range:
assert r1.scenario_range.rates == r2.scenario_range.rates
+54 -3
View File
@@ -1,7 +1,9 @@
"""Raw sample pruning tests.
Spec §3.4, ST-5: Raw samples pruned to 14 days; hour observations and
day aggregates retained indefinitely.
day aggregates are retained indefinitely.
Issue #74: Boundary anchors required for successor evidence are retained.
"""
import sqlite3
import sys
@@ -51,6 +53,9 @@ class TestPruneOldSamples:
def test_removes_old_samples(self, store_conn):
now = datetime(2026, 9, 15, 12, 0, 0, tzinfo=timezone.utc)
# Insert samples at 10, 14, and 15 days ago
# All three are before the cutoff (2026-09-01T12:00:00)
# The 15-day-old sample is not a boundary anchor because
# the next sample (14 days ago) is also before the cutoff
for days_ago in [10, 14, 15]:
ts = (now - timedelta(days=days_ago)).isoformat()
_insert_sample(store_conn, ts)
@@ -82,6 +87,7 @@ class TestPruneOldSamples:
"""Hour observations are retained indefinitely."""
now = datetime(2026, 9, 15, 12, 0, 0, tzinfo=timezone.utc)
# Insert an old sample and a recent sample
# The old sample is a boundary anchor (needed for derivation)
_insert_sample(store_conn, (now - timedelta(days=20)).isoformat())
_insert_sample(store_conn, (now - timedelta(days=1)).isoformat())
@@ -95,10 +101,55 @@ class TestPruneOldSamples:
prune_old_samples(store_conn, now, retention_days=14)
# Sample removed
# Old sample is retained as boundary anchor (needed for derivation)
cursor = store_conn.execute("SELECT COUNT(*) FROM samples")
assert cursor.fetchone()[0] == 1
assert cursor.fetchone()[0] == 2
# Hour observation retained
cursor = store_conn.execute("SELECT COUNT(*) FROM hour_observations")
assert cursor.fetchone()[0] == 1
def test_boundary_anchor_retained(self, store_conn):
"""Boundary anchors required for derivation are retained."""
now = datetime(2026, 9, 30, 12, 0, 0, tzinfo=timezone.utc)
# Old sample before boundary (2026-09-14T23:55:00)
# is 15 days and 0.75 hours old (before cutoff at 2026-09-16T12:00:00)
_insert_sample(store_conn, "2026-09-14T23:55:00+00:00", 1000000)
# Sample after boundary (2026-09-16T12:30:00)
# is 13 days and 23.5 hours old (within retention)
_insert_sample(store_conn, "2026-09-16T12:30:00+00:00", 2000000)
# Run pruning
pruned = prune_old_samples(store_conn, now, retention_days=14)
# The boundary anchor should be retained
cursor = store_conn.execute(
"SELECT COUNT(*) FROM samples WHERE ts = '2026-09-14T23:55:00+00:00'"
)
assert cursor.fetchone()[0] == 1
def test_old_sample_with_derived_interval_removed(self, store_conn):
"""Old samples with fully derived intervals are removed."""
now = datetime(2026, 9, 30, 12, 0, 0, tzinfo=timezone.utc)
# Old sample with derived interval
_insert_sample(store_conn, "2026-09-10T10:00:00+00:00", 1000000)
_insert_sample(store_conn, "2026-09-10T10:30:00+00:00", 2000000)
# Hour observation exists for the interval
store_conn.execute(
"INSERT INTO hour_observations (hour, active_seconds, bytes_written_delta, bytes_read_delta, sample_count, coverage) "
"VALUES ('2026-09-10T10:00:00+00:00', 3600, 1000000, 0, 1, 1.0)",
)
store_conn.commit()
# Run pruning
pruned = prune_old_samples(store_conn, now, retention_days=14)
# Old sample should be removed (interval is derived)
cursor = store_conn.execute(
"SELECT COUNT(*) FROM samples WHERE ts = '2026-09-10T10:00:00+00:00'"
)
assert cursor.fetchone()[0] == 0
+398
View File
@@ -0,0 +1,398 @@
"""Repair and retention tests (issue #74).
Tests the idempotent, safe repair of hour observations and day aggregates
from surviving raw samples, boundary anchor retention, and legacy summary
handling at actual precision.
"""
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.store import init_store
from fenris.repair import (
repair_derivation,
is_repair_in_progress,
get_repair_status,
)
from fenris.pruning import prune_old_samples, needs_boundary_anchor
# ---------------------------------------------------------------------------
# Fixtures
# ---------------------------------------------------------------------------
@pytest.fixture
def store_conn(tmp_path: Path):
"""Create a fresh store for each test."""
db_path = tmp_path / "test.db"
conn = init_store(db_path)
yield conn
conn.close()
def _insert_sample(conn, ts_iso, bytes_written, bytes_read=0, power_on_hours=100,
device="/dev/nvme0", segment_id=None):
"""Insert a raw sample."""
conn.execute(
"""INSERT INTO samples
(ts, device, bytes_written, bytes_read, power_on_hours,
data_units_written, data_units_read, segment_id)
VALUES (?, ?, ?, ?, ?, ?, ?, ?)""",
(ts_iso, device, bytes_written, bytes_read, power_on_hours,
bytes_written // 512000, bytes_read // 512000, segment_id),
)
conn.commit()
def _insert_hour(conn, hour_iso, bytes_written_delta=0, active_seconds=3600,
idle_seconds=0, powered_off_seconds=0, unknown_seconds=0,
sample_count=1, coverage=1.0):
"""Insert an hour observation."""
conn.execute(
"""INSERT INTO hour_observations
(hour, active_seconds, idle_seconds, powered_off_seconds, unknown_seconds,
bytes_written_delta, bytes_read_delta, sample_count, coverage)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)""",
(hour_iso, active_seconds, idle_seconds, powered_off_seconds, unknown_seconds,
bytes_written_delta, 0, sample_count, coverage),
)
conn.commit()
def _insert_day_aggregate(conn, day, bytes_written_delta=0, coverage=1.0):
"""Insert a day aggregate."""
conn.execute(
"""INSERT INTO day_aggregates
(day, active_seconds, bytes_written_delta, coverage, sample_count)
VALUES (?, 3600, ?, ?, 1)""",
(day, bytes_written_delta, coverage),
)
conn.commit()
def _open_period(conn, start_iso, end_iso=None, end_cause=None):
"""Insert a monitoring period."""
conn.execute(
"""INSERT INTO monitoring_periods (started_at, ended_at, end_cause)
VALUES (?, ?, ?)""",
(start_iso, end_iso, end_cause),
)
conn.commit()
# ---------------------------------------------------------------------------
# AC1: Normal collection, legacy import and recovery use same evidence rules
# ---------------------------------------------------------------------------
class TestRepairUsesSameEvidenceRules:
"""AC1: Repair derives only surviving supported evidence, transactionally
and idempotently."""
def test_repair_idempotent_on_empty_store(self, store_conn):
"""Repair on empty store succeeds and does nothing."""
result = repair_derivation(store_conn)
assert result.ok is True
assert result.hours_created == 0
assert result.days_created == 0
def test_repair_idempotent_on_fully_derived(self, store_conn):
"""Repair on store with existing derived data does not duplicate."""
now = datetime(2026, 9, 15, 12, 0, 0, tzinfo=timezone.utc)
_open_period(store_conn, "2026-09-14T00:00:00+00:00")
# Insert samples that span an hour
_insert_sample(store_conn, "2026-09-14T10:00:00+00:00", 1000000)
_insert_sample(store_conn, "2026-09-14T10:30:00+00:00", 2000000)
# Pre-existing hour observation for the 10:00 hour
# This hour observation captures the interval [10:00, 10:30]
_insert_hour(store_conn, "2026-09-14T10:00:00+00:00",
bytes_written_delta=1000000)
_insert_day_aggregate(store_conn, "2026-09-14",
bytes_written_delta=1000000)
# Run repair
result = repair_derivation(store_conn)
# Should not create new observations (already exist)
assert result.hours_created == 0
assert result.days_created == 0
# Verify no duplicates
cursor = store_conn.execute(
"SELECT COUNT(*) FROM hour_observations WHERE hour = '2026-09-14T10:00:00+00:00'"
)
assert cursor.fetchone()[0] == 1
def test_repair_preserves_import_markers(self, store_conn):
"""Repair does not remove legacy import markers."""
# Set legacy import marker
store_conn.execute(
"INSERT INTO store_metadata (key, value) VALUES ('legacy_imported', 'true')"
)
store_conn.commit()
# Run repair
result = repair_derivation(store_conn)
# Verify marker preserved
cursor = store_conn.execute(
"SELECT value FROM store_metadata WHERE key = 'legacy_imported'"
)
assert cursor.fetchone()[0] == "true"
# ---------------------------------------------------------------------------
# AC2: Rerunning or interrupting repair produces no duplicates
# ---------------------------------------------------------------------------
class TestRepairIdempotency:
"""AC2: No duplicated intervals or totals from repeated repair."""
def test_repair_no_duplicate_hours(self, store_conn):
"""Running repair twice produces no duplicate hour observations."""
_open_period(store_conn, "2026-09-14T00:00:00+00:00")
_insert_sample(store_conn, "2026-09-14T10:00:00+00:00", 1000000)
_insert_sample(store_conn, "2026-09-14T10:30:00+00:00", 2000000)
# First repair
result1 = repair_derivation(store_conn)
assert result1.hours_created == 1
# Second repair
result2 = repair_derivation(store_conn)
assert result2.hours_created == 0 # No new hours
# Verify only one hour observation
cursor = store_conn.execute("SELECT COUNT(*) FROM hour_observations")
assert cursor.fetchone()[0] == 1
def test_repair_no_duplicate_days(self, store_conn):
"""Running repair twice produces no duplicate day aggregates."""
_open_period(store_conn, "2026-09-14T00:00:00+00:00")
_insert_sample(store_conn, "2026-09-14T10:00:00+00:00", 1000000)
_insert_sample(store_conn, "2026-09-14T10:30:00+00:00", 2000000)
# First repair
result1 = repair_derivation(store_conn)
assert result1.days_created == 1
# Second repair
result2 = repair_derivation(store_conn)
assert result2.days_created == 0 # No new days
# Verify only one day aggregate
cursor = store_conn.execute("SELECT COUNT(*) FROM day_aggregates")
assert cursor.fetchone()[0] == 1
def test_interrupted_repair_preserves_evidence(self, store_conn):
"""If repair fails, prior valid history is preserved."""
_open_period(store_conn, "2026-09-14T00:00:00+00:00")
# Insert valid existing data
_insert_hour(store_conn, "2026-09-14T10:00:00+00:00",
bytes_written_delta=500000)
_insert_day_aggregate(store_conn, "2026-09-14",
bytes_written_delta=500000)
# Run repair (should succeed but not modify existing valid data)
result = repair_derivation(store_conn)
assert result.ok is True
# Verify existing data preserved
cursor = store_conn.execute(
"SELECT bytes_written_delta FROM hour_observations "
"WHERE hour = '2026-09-14T10:00:00+00:00'"
)
assert cursor.fetchone()[0] == 500000
# ---------------------------------------------------------------------------
# AC3: Boundary anchor retention
# ---------------------------------------------------------------------------
class TestBoundaryAnchorRetention:
"""AC3: Prune samples only after durable derivation; retain anchors."""
def test_needs_boundary_anchor_sample(self, store_conn):
"""Sample before 14-day boundary is needed for derivation."""
now = datetime(2026, 9, 30, 12, 0, 0, tzinfo=timezone.utc)
# Sample just before 14-day boundary (2026-09-15T23:55:00)
# is 14 days and 0.75 hours old (before cutoff at 2026-09-16T12:00:00)
_insert_sample(store_conn, "2026-09-15T23:55:00+00:00", 1000000)
# Sample just after boundary (2026-09-16T12:30:00)
# is 13 days and 23.5 hours old (within retention)
_insert_sample(store_conn, "2026-09-16T12:30:00+00:00", 2000000)
# The sample at 2026-09-15 is a boundary anchor because
# the interval spans the retention boundary
assert needs_boundary_anchor(store_conn, "2026-09-15T23:55:00+00:00", now)
def test_not_boundary_anchor_if_fully_derived(self, store_conn):
"""Sample that's fully derived is not a boundary anchor."""
now = datetime(2026, 9, 30, 12, 0, 0, tzinfo=timezone.utc)
# Insert sample and fully derive its interval
_insert_sample(store_conn, "2026-09-14T10:00:00+00:00", 1000000)
_insert_sample(store_conn, "2026-09-14T10:30:00+00:00", 2000000)
# Hour observation already exists for this interval
_insert_hour(store_conn, "2026-09-14T10:00:00+00:00",
bytes_written_delta=1000000)
# Not a boundary anchor
assert not needs_boundary_anchor(store_conn, "2026-09-14T10:00:00+00:00", now)
def test_pruning_retains_boundary_anchors(self, store_conn):
"""Pruning keeps samples needed as boundary anchors."""
now = datetime(2026, 9, 30, 12, 0, 0, tzinfo=timezone.utc)
# Old sample before boundary (2026-09-14T23:55:00)
# is 15 days and 0.75 hours old (before cutoff at 2026-09-16T12:00:00)
_insert_sample(store_conn, "2026-09-14T23:55:00+00:00", 1000000)
# Sample after boundary (2026-09-16T12:30:00)
# is 13 days and 23.5 hours old (within retention)
_insert_sample(store_conn, "2026-09-16T12:30:00+00:00", 2000000)
# Recent sample
_insert_sample(store_conn, "2026-09-29T12:00:00+00:00", 3000000)
# Run pruning
pruned = prune_old_samples(store_conn, now, retention_days=14)
# The boundary anchor should be retained
cursor = store_conn.execute(
"SELECT COUNT(*) FROM samples WHERE ts = '2026-09-14T23:55:00+00:00'"
)
assert cursor.fetchone()[0] == 1
def test_pruning_removes_old_sample_with_derived_interval(self, store_conn):
"""Pruning removes old samples when interval is fully derived."""
now = datetime(2026, 9, 30, 12, 0, 0, tzinfo=timezone.utc)
# Old sample with derived interval
_insert_sample(store_conn, "2026-09-10T10:00:00+00:00", 1000000)
_insert_sample(store_conn, "2026-09-10T10:30:00+00:00", 2000000)
# Hour observation exists for the interval
_insert_hour(store_conn, "2026-09-10T10:00:00+00:00",
bytes_written_delta=1000000)
# Run pruning
pruned = prune_old_samples(store_conn, now, retention_days=14)
# Old sample should be removed (interval is derived)
cursor = store_conn.execute(
"SELECT COUNT(*) FROM samples WHERE ts = '2026-09-10T10:00:00+00:00'"
)
assert cursor.fetchone()[0] == 0
# ---------------------------------------------------------------------------
# AC4: Legacy day-only summaries retain actual precision
# ---------------------------------------------------------------------------
class TestLegacySummaryPrecision:
"""AC4: Legacy summaries at actual precision, no interpolation."""
def test_legacy_summary_not_reconstructed(self, store_conn):
"""Legacy day-only summaries are not interpolated to hour detail."""
# Insert a legacy-style day aggregate without hour observations
_insert_day_aggregate(store_conn, "2026-08-01",
bytes_written_delta=5000000)
# Run repair
result = repair_derivation(store_conn)
# Should not create hour observations for legacy day
cursor = store_conn.execute(
"SELECT COUNT(*) FROM hour_observations WHERE hour LIKE '2026-08-01%'"
)
assert cursor.fetchone()[0] == 0
def test_legacy_summary_no_double_counting(self, store_conn):
"""Legacy summaries and derived intervals don't double-count."""
# Insert legacy day aggregate
_insert_day_aggregate(store_conn, "2026-08-01",
bytes_written_delta=5000000)
# Run repair
result = repair_derivation(store_conn)
# Day aggregate should not be modified
cursor = store_conn.execute(
"SELECT bytes_written_delta FROM day_aggregates WHERE day = '2026-08-01'"
)
assert cursor.fetchone()[0] == 5000000
# ---------------------------------------------------------------------------
# AC5: Status distinguishes evidence states
# ---------------------------------------------------------------------------
class TestStatusEvidenceDistingushing:
"""AC5: Read-only views distinguish evidence states."""
def test_repair_status_available(self, store_conn):
"""Repair status is available for read-only views."""
status = get_repair_status(store_conn)
assert hasattr(status, 'last_repair')
assert hasattr(status, 'repair_in_progress')
assert hasattr(status, 'hours_derived')
assert hasattr(status, 'days_derived')
def test_repair_in_progress_flag(self, store_conn):
"""Repair in progress flag is trackable."""
assert not is_repair_in_progress(store_conn)
# ---------------------------------------------------------------------------
# AC6: Migration then collection then reader consumption
# ---------------------------------------------------------------------------
class TestMigrationCollectionReader:
"""AC6: End-to-end migration, collection, and reader consumption."""
def test_store_with_raw_evidence_and_summaries(self, store_conn):
"""Store with raw evidence and old summaries works correctly."""
# Set up store with mixed data
_open_period(store_conn, "2026-08-01T00:00:00+00:00")
# Old day-only summary (legacy)
_insert_day_aggregate(store_conn, "2026-08-01",
bytes_written_delta=5000000)
# Recent raw samples
_insert_sample(store_conn, "2026-09-14T10:00:00+00:00", 1000000)
_insert_sample(store_conn, "2026-09-14T10:30:00+00:00", 2000000)
# Run repair
result = repair_derivation(store_conn)
assert result.ok is True
# Verify legacy summary preserved
cursor = store_conn.execute(
"SELECT bytes_written_delta FROM day_aggregates WHERE day = '2026-08-01'"
)
assert cursor.fetchone()[0] == 5000000
# Verify new hour observation created
cursor = store_conn.execute(
"SELECT COUNT(*) FROM hour_observations WHERE hour LIKE '2026-09-14%'"
)
assert cursor.fetchone()[0] == 1
def test_concurrent_reader_consistency(self, store_conn):
"""Reader sees consistent snapshot during repair."""
# This is more of a documentation test - SQLite WAL mode handles this
# We verify the store is in WAL mode
cursor = store_conn.execute("PRAGMA journal_mode")
assert cursor.fetchone()[0] == "wal"
+4 -4
View File
@@ -118,12 +118,12 @@ class TestMigrateToLatest:
assert migrate_to_latest(db) == 0
def test_migrates_intermediate_version(self, tmp_path):
"""Store at version 1 with SCHEMA_VERSION=1 → 0 steps (current)."""
"""Store at version SCHEMA_VERSION-1 → 1 step to current."""
from fenris.store import SCHEMA_VERSION
db = tmp_path / "observations.db"
_make_store(db, version=1)
# SCHEMA_VERSION is 1, so version 1 is current
_make_store(db, version=SCHEMA_VERSION - 1)
steps = migrate_to_latest(db)
assert steps == 0
assert steps == 1
# ---------------------------------------------------------------------------
+407
View File
@@ -40,12 +40,23 @@ from fenris.status import (
)
from fenris.tui import (
FenrisTuiApp,
DailyBarGraph,
_format_remaining,
_sparkline,
_habit_bar,
_query_usage_history,
_query_drive_health,
_query_service_facts,
_query_daily_graph_data,
_query_hourly_graph_data,
_RANGE_OPTIONS,
_RANGE_DEFAULT,
_BAR_HEIGHT,
_GLYPH_ALLOCATED,
_GLYPH_UNALLOCATED,
_GLYPH_GAP,
_GLYPH_ZERO,
_GLYPH_PARTIAL,
)
@@ -679,3 +690,399 @@ class TestStateMatrixCombinations:
# Incomplete provenance → UNVERIFIED tier
assert proj.baseline_tier.value == "unverified_override"
conn.close()
# ---------------------------------------------------------------------------
# Hour observation helper for graph tests
# ---------------------------------------------------------------------------
def _insert_hour(conn, hour, bw=0, active=3600, idle=0, powered_off=0,
unknown=0, coverage=1.0, samples=1):
"""Insert an hour observation row."""
conn.execute(
"INSERT INTO hour_observations "
"(hour, active_seconds, idle_seconds, powered_off_seconds, unknown_seconds, "
" bytes_written_delta, bytes_read_delta, sample_count, coverage) "
"VALUES (?, ?, ?, ?, ?, ?, 0, ?, ?)",
(hour, active, idle, powered_off, unknown, bw, samples, coverage),
)
conn.commit()
def _insert_day_with_unattributed(conn, day, bw=0, unattributed=0,
coverage=0.95, samples=24):
"""Insert a day aggregate with explicit unattributed bytes."""
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, "
" unattributed_bytes_written, unattributed_bytes_read) "
"VALUES (?, 3600, 0, 0, 0, ?, 0, ?, ?, ?, 0)",
(day, bw, samples, coverage, unattributed),
)
conn.commit()
# ---------------------------------------------------------------------------
# Graph data query tests (issue #75)
# ---------------------------------------------------------------------------
class TestQueryDailyGraphData:
def test_empty_store(self, tmp_path):
conn = init_store(tmp_path / "test.db")
result = _query_daily_graph_data(conn)
assert result == []
conn.close()
def test_with_days(self, tmp_path):
conn = init_store(tmp_path / "test.db")
_insert_segment(conn)
_open_period(conn)
for i in range(14):
d = (datetime(2026, 9, 15) + timedelta(days=i)).strftime("%Y-%m-%d")
_insert_day(conn, d, bw=1024*1024*100)
result = _query_daily_graph_data(conn)
assert len(result) == 14
assert result[0]["total_bytes"] == 1024*1024*100
assert result[0]["allocated_bytes"] == 1024*1024*100
assert result[0]["unallocated_bytes"] == 0
assert result[0]["is_zero"] is False
conn.close()
def test_zero_day(self, tmp_path):
conn = init_store(tmp_path / "test.db")
_insert_segment(conn)
_open_period(conn)
_insert_day(conn, "2026-09-20", bw=0, coverage=0.95, samples=24)
result = _query_daily_graph_data(conn)
assert len(result) == 1
assert result[0]["is_zero"] is True
assert result[0]["total_bytes"] == 0
conn.close()
def test_unattributed_bytes(self, tmp_path):
conn = init_store(tmp_path / "test.db")
_insert_segment(conn)
_open_period(conn)
_insert_day_with_unattributed(
conn, "2026-09-20", bw=500, unattributed=300
)
result = _query_daily_graph_data(conn)
assert len(result) == 1
assert result[0]["allocated_bytes"] == 500
assert result[0]["unallocated_bytes"] == 300
assert result[0]["total_bytes"] == 800
conn.close()
def test_gap_day(self, tmp_path):
conn = init_store(tmp_path / "test.db")
_insert_segment(conn)
_open_period(conn)
# Day with no hours but unknown seconds (gap in monitoring period)
conn.execute(
"INSERT INTO day_aggregates "
"(day, active_seconds, idle_seconds, powered_off_seconds, unknown_seconds, "
" bytes_written_delta, sample_count, coverage) "
"VALUES (?, 0, 0, 0, 86400, 0, 0, 0.0)",
("2026-09-20",),
)
conn.commit()
result = _query_daily_graph_data(conn)
assert len(result) == 1
assert result[0]["is_gap"] is True
conn.close()
class TestQueryHourlyGraphData:
def test_empty_day(self, tmp_path):
conn = init_store(tmp_path / "test.db")
result = _query_hourly_graph_data(conn, "2026-09-20")
assert result == []
conn.close()
def test_with_hours(self, tmp_path):
conn = init_store(tmp_path / "test.db")
for h in range(24):
hour = "2026-09-20T%02d:00:00+00:00" % h
bw = 1024*1024*100 if h in (10, 14) else 0
_insert_hour(conn, hour, bw=bw)
result = _query_hourly_graph_data(conn, "2026-09-20")
assert len(result) == 24
assert result[10]["bytes_written"] == 1024*1024*100
assert result[10]["is_zero"] is False
assert result[0]["is_zero"] is True
conn.close()
def test_hour_labels(self, tmp_path):
conn = init_store(tmp_path / "test.db")
_insert_hour(conn, "2026-09-20T12:00:00+00:00", bw=1000)
result = _query_hourly_graph_data(conn, "2026-09-20")
assert len(result) == 1
assert result[0]["local_label"] == "12"
conn.close()
# ---------------------------------------------------------------------------
# DailyBarGraph widget unit tests (issue #75)
# ---------------------------------------------------------------------------
class TestDailyBarGraph:
def test_empty_data(self, tmp_path):
"""Empty data shows awaiting message."""
app = FenrisTuiApp(store_path=tmp_path / "nonexistent.db")
# We test the widget directly via the app's compose
graph = DailyBarGraph()
# Simulate setting empty data
from textual.app import App as TextualApp
class _TestApp(TextualApp):
def compose(self):
yield graph
# We can't easily test widget lifecycle outside an app, so test the data
assert graph._all_day_data == []
assert graph.range_days == _RANGE_DEFAULT
def test_range_default(self):
"""Default range is 14 days."""
graph = DailyBarGraph()
assert graph.range_days == _RANGE_DEFAULT
assert graph.selected_index == -1
assert graph.view_mode == "daily"
def test_trim_to_range(self):
"""Data is trimmed to the selected range."""
graph = DailyBarGraph()
data = [{"day": "2026-09-%02d" % d, "total_bytes": d * 100}
for d in range(1, 31)]
graph._all_day_data = data
graph.range_days = 7
graph._trim_to_range()
assert len(graph._day_data) == 7
assert graph._day_data[0]["day"] == "2026-09-24"
def test_range_all_fits(self):
"""When data fits within range, all days are shown."""
graph = DailyBarGraph()
data = [{"day": "2026-09-%02d" % d, "total_bytes": d * 100}
for d in range(1, 8)]
graph._all_day_data = data
graph.range_days = 14
graph._trim_to_range()
assert len(graph._day_data) == 7
# ---------------------------------------------------------------------------
# Bar graph headless TUI tests (issue #75)
# ---------------------------------------------------------------------------
class TestBarGraphTUI:
"""Headless tests for the interactive bar graph."""
@pytest.mark.asyncio
async def test_graph_renders_with_data(self, tmp_path):
"""Graph widget renders with day data."""
conn = init_store(tmp_path / "test.db")
_insert_segment(conn)
_open_period(conn)
for i in range(14):
d = (datetime(2026, 9, 15) + timedelta(days=i)).strftime("%Y-%m-%d")
_insert_day(conn, d, bw=1024*1024*100)
_insert_sample(conn, "2026-09-30T10:00:00+00:00")
conn.close()
app = FenrisTuiApp(store_path=tmp_path / "test.db")
async with app.run_test(size=(80, 24)) as pilot:
graph = app.query_one("#usage-history")
assert isinstance(graph, DailyBarGraph)
assert len(graph._day_data) > 0
# Legend should be visible
legend = str(app.query_one("#bar-legend").render())
assert "alloc" in legend
assert "zero" in legend
@pytest.mark.asyncio
async def test_arrow_selection(self, tmp_path):
"""Left/right arrows move the selection."""
conn = init_store(tmp_path / "test.db")
_insert_segment(conn)
_open_period(conn)
for i in range(14):
d = (datetime(2026, 9, 15) + timedelta(days=i)).strftime("%Y-%m-%d")
_insert_day(conn, d, bw=1024*1024*100)
_insert_sample(conn, "2026-09-30T10:00:00+00:00")
conn.close()
app = FenrisTuiApp(store_path=tmp_path / "test.db")
async with app.run_test(size=(80, 24)) as pilot:
graph = app.query_one("#usage-history")
app.set_focus(graph)
await pilot.pause()
await pilot.pause()
# Right arrow selects first bar
await pilot.press("right")
await pilot.pause()
assert graph.selected_index == 0
# Right again moves to second
await pilot.press("right")
await pilot.pause()
assert graph.selected_index == 1
# Left moves back
await pilot.press("left")
await pilot.pause()
assert graph.selected_index == 0
@pytest.mark.asyncio
async def test_range_switching(self, tmp_path):
"""1-4 keys switch the visible range."""
conn = init_store(tmp_path / "test.db")
_insert_segment(conn)
_open_period(conn)
for i in range(30):
d = (datetime(2026, 9, 1) + timedelta(days=i)).strftime("%Y-%m-%d")
_insert_day(conn, d, bw=1024*1024*100)
_insert_sample(conn, "2026-09-30T10:00:00+00:00")
conn.close()
app = FenrisTuiApp(store_path=tmp_path / "test.db")
async with app.run_test(size=(80, 24)) as pilot:
graph = app.query_one("#usage-history")
app.set_focus(graph)
await pilot.pause()
await pilot.pause()
# Default is 14
assert graph.range_days == 14
assert len(graph._day_data) == 14
# Press 1 for 7-day range
await pilot.press("1")
await pilot.pause()
assert graph.range_days == 7
assert len(graph._day_data) == 7
# Press 3 for 28-day range
await pilot.press("3")
await pilot.pause()
assert graph.range_days == 28
assert len(graph._day_data) == 28
# Press 4 for 90-day range (only 30 days available)
await pilot.press("4")
await pilot.pause()
assert graph.range_days == 90
assert len(graph._day_data) == 30
@pytest.mark.asyncio
async def test_readout_updates_on_selection(self, tmp_path):
"""Readout shows selected day info."""
conn = init_store(tmp_path / "test.db")
_insert_segment(conn)
_open_period(conn)
for i in range(14):
d = (datetime(2026, 9, 15) + timedelta(days=i)).strftime("%Y-%m-%d")
_insert_day(conn, d, bw=1024*1024*100)
_insert_sample(conn, "2026-09-30T10:00:00+00:00")
conn.close()
app = FenrisTuiApp(store_path=tmp_path / "test.db")
async with app.run_test(size=(80, 24)) as pilot:
graph = app.query_one("#usage-history")
app.set_focus(graph)
await pilot.pause()
await pilot.pause()
# No selection initially
readout = str(app.query_one("#bar-readout").render())
assert "select" in readout.lower()
# Select first bar
await pilot.press("right")
await pilot.pause()
readout = str(app.query_one("#bar-readout").render())
assert "2026-09" in readout
assert "GB" in readout
@pytest.mark.asyncio
async def test_hourly_drill_down_and_back(self, tmp_path):
"""Enter drills into hourly view, Esc returns to daily."""
conn = init_store(tmp_path / "test.db")
_insert_segment(conn)
_open_period(conn)
for i in range(14):
d = (datetime(2026, 9, 15) + timedelta(days=i)).strftime("%Y-%m-%d")
_insert_day(conn, d, bw=1024*1024*100)
# Insert hours for each day
for h in range(24):
hour = "%sT%02d:00:00+00:00" % (d, h)
bw_h = 1024*1024*10 if h == 12 else 0
_insert_hour(conn, hour, bw=bw_h)
_insert_sample(conn, "2026-09-30T10:00:00+00:00")
conn.close()
app = FenrisTuiApp(store_path=tmp_path / "test.db")
async with app.run_test(size=(80, 24)) as pilot:
graph = app.query_one("#usage-history")
app.set_focus(graph)
await pilot.pause()
await pilot.pause()
# Select a day
await pilot.press("right")
await pilot.pause()
assert graph.selected_index == 0
assert graph.view_mode == "daily"
# Enter drill-down
await pilot.press("enter")
await pilot.pause()
assert graph.view_mode == "hourly"
assert graph.drill_day is not None
assert len(graph._hour_data) == 24
# Esc returns to daily
await pilot.press("escape")
await pilot.pause()
assert graph.view_mode == "daily"
assert graph.drill_day is None
@pytest.mark.asyncio
async def test_empty_store_graph(self, tmp_path):
"""Empty store shows awaiting message in graph."""
app = FenrisTuiApp(store_path=tmp_path / "nonexistent.db")
async with app.run_test(size=(80, 24)) as pilot:
graph = app.query_one("#usage-history")
assert isinstance(graph, DailyBarGraph)
render = str(app.query_one("#bar-render").render())
assert "awaiting" in render.lower()
@pytest.mark.asyncio
async def test_graph_border_title(self, tmp_path):
"""Graph pane has a border title."""
app = FenrisTuiApp(store_path=tmp_path / "nonexistent.db")
async with app.run_test(size=(80, 24)) as pilot:
graph = app.query_one("#usage-history")
assert graph.border_title == "usage history"
@pytest.mark.asyncio
async def test_glyphs_in_legend(self, tmp_path):
"""Legend shows all required glyphs."""
conn = init_store(tmp_path / "test.db")
_insert_segment(conn)
_open_period(conn)
for i in range(7):
d = (datetime(2026, 9, 20) + timedelta(days=i)).strftime("%Y-%m-%d")
_insert_day(conn, d, bw=1024*1024*100)
_insert_sample(conn, "2026-09-30T10:00:00+00:00")
conn.close()
app = FenrisTuiApp(store_path=tmp_path / "test.db")
async with app.run_test(size=(80, 24)) as pilot:
await pilot.pause()
legend = str(app.query_one("#bar-legend").render())
assert _GLYPH_ALLOCATED in legend
assert _GLYPH_UNALLOCATED in legend
assert _GLYPH_GAP in legend
assert _GLYPH_ZERO in legend
assert _GLYPH_PARTIAL in legend