Compare commits
2
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
6917a658cb | ||
|
|
d790ff84c5 |
@@ -1,67 +0,0 @@
|
||||
# PROTOTYPE — titlebox, drive health and colour presets (throwaway)
|
||||
|
||||
Answers wayfinder ticket **Approve titlebox, health layout and colour presets**
|
||||
on map **Fenris TUI polish and hourly history**.
|
||||
Not production code. Do not merge onto main.
|
||||
|
||||
## Question
|
||||
|
||||
What concrete layout and preset contract makes Fenris feel polished while
|
||||
preserving monitoring clarity and accessibility?
|
||||
|
||||
## Run
|
||||
|
||||
```sh
|
||||
./run
|
||||
# then open http://127.0.0.1:8765/prototype/titlebox-health-presets/
|
||||
```
|
||||
|
||||
Or open `index.html` directly.
|
||||
|
||||
## What this prototype asks the human to react to
|
||||
|
||||
Use the bottom control bar to switch:
|
||||
|
||||
- preset: **Amber**, **Nord**, **High Contrast**;
|
||||
- status: Monitoring, Collecting, Paused, Waiting, Interrupted, Error, Stale,
|
||||
Unknown;
|
||||
- terminal size: 80×24 baseline, 140×40 wide, or a constrained/narrow layout;
|
||||
- title glyph: wolf emoji or fallback text;
|
||||
- motion: normal blink or reduced motion.
|
||||
|
||||
## Contract under test
|
||||
|
||||
1. **Title hierarchy** — the only maker-credit surface is the top titlebox,
|
||||
`🐺 Fenris by Bongbetic`; the old service-strip credit is removed. The
|
||||
lifespan headline stays a product data surface, not the app title.
|
||||
2. **Glyph fallback** — when the wolf is unsupported, explicitly fall back to
|
||||
`Fenris by Bongbetic` rather than rendering tofu or disturbing border width.
|
||||
3. **Drive health vs settings** — vendor wear lives in **Drive health** beside
|
||||
temperature/spare/error facts. **Settings** stays read-only for device,
|
||||
baseline, retention and display preferences.
|
||||
4. **Palette roles** — themes style chrome, graph and accent colours; status
|
||||
semantics (green/amber/red/grey plus glyph+text) outrank the theme and are
|
||||
never the sole carrier.
|
||||
5. **Display preferences** — preset and reduced-motion preferences are scoped
|
||||
to the unprivileged TUI user, not `/etc/fenris/fenris.conf` and not the
|
||||
observation store. Proposed persistence: `${XDG_CONFIG_HOME:-~/.config}/fenris/tui.json`.
|
||||
6. **Reduced motion** — normal Monitoring blinks only the dot at the settled
|
||||
750 ms cadence; reduced motion renders the same `● Monitoring` state steady.
|
||||
7. **Constrained terminal** — below 80×24, the graph region hides and leaves a
|
||||
one-line textual summary plus `graph needs ≥80×24`; status, title, drive
|
||||
health and action/quit affordances stay visible.
|
||||
|
||||
## Framework facts verified
|
||||
|
||||
Current Textual docs via Context7 library `/websites/textual_textualize_io`:
|
||||
|
||||
- `App.register_theme(theme)` registers a theme; setting `App.theme` activates
|
||||
it. Textual themes generate CSS variables from base colours.
|
||||
- `$text`, `$text-muted`, `$text-disabled` and `color: auto` exist to maintain
|
||||
text legibility against theme backgrounds.
|
||||
- Mouse events expose coordinates relative to the screen or widget; focusable
|
||||
widgets can be resolved at coordinates, and Pilot can click widgets/offsets
|
||||
for automated proof of keyboard+mouse paths.
|
||||
|
||||
Local repo lockfile still pins `textual==8.2.8`; this prototype is static HTML
|
||||
so it deliberately does not depend on the Python environment.
|
||||
@@ -1,89 +0,0 @@
|
||||
<!doctype html>
|
||||
<html lang="en">
|
||||
<head>
|
||||
<meta charset="utf-8">
|
||||
<meta name="viewport" content="width=device-width,initial-scale=1">
|
||||
<title>Fenris titlebox, health layout and presets prototype</title>
|
||||
<style>
|
||||
:root{color-scheme:dark;font-family:ui-monospace,SFMono-Regular,Consolas,Menlo,monospace}
|
||||
*{box-sizing:border-box}body{margin:0;background:#06070a;color:var(--text);font:14px/1.35 ui-monospace,SFMono-Regular,Consolas,Menlo,monospace}.note{max-width:1180px;margin:18px auto 8px;padding:0 16px;color:#c0c6d0}.terminal{width:min(var(--term-w),calc(100vw - 32px));min-height:var(--term-h);margin:0 auto 108px;background:var(--bg);border:1px solid var(--line);box-shadow:0 24px 80px #000;border-radius:8px;overflow:hidden}.screen{padding:8px}.titlebox{height:38px;border:2px solid var(--accent);display:flex;align-items:center;justify-content:center;font-weight:900;font-size:17px;letter-spacing:.01em;color:var(--title);background:linear-gradient(180deg,var(--title-bg),transparent)}.titlebox .fallback{display:none}body[data-glyph=plain] .titlebox .emoji{display:none}body[data-glyph=plain] .titlebox .fallback{display:inline}.auth{margin-top:6px;color:var(--muted);border-left:4px solid var(--accent);padding:4px 8px;background:var(--panel-soft)}.grid{display:grid;grid-template-columns:3fr 2fr;gap:8px;margin-top:8px}.pane{border:1px solid var(--line);background:var(--panel);padding:9px 11px;min-height:130px}.headline{grid-column:1/-1;min-height:92px}.pane-title{color:var(--muted);margin:-19px 0 6px -2px;background:var(--panel);width:max-content;padding:0 6px}.big{font-size:18px;font-weight:900}.muted{color:var(--muted)}.accent{color:var(--accent)}.status-line{display:flex;gap:12px;align-items:center;justify-content:space-between;flex-wrap:wrap;margin-top:5px}.state{font-weight:900}.state .glyph{display:inline-block;width:2ch;text-align:center}.state.monitoring{color:var(--ok)}body:not(.reduced) .state.monitoring .glyph{animation:blink 1.5s steps(1,end) infinite}.state.collecting{color:var(--ok)}.state.paused,.state.waiting{color:var(--warn)}.state.interrupted,.state.error,.state.stale{color:var(--bad)}.state.unknown{color:var(--neutral)}@keyframes blink{50%{opacity:.18}}.pause-banner{display:none;margin-top:8px;border:2px solid var(--warn);background:var(--warn-bg);color:var(--warn-text);padding:7px 10px;font-weight:900}body[data-state=paused] .pause-banner{display:block}.facts{display:grid;grid-template-columns:repeat(4,minmax(0,1fr));gap:5px;margin-top:7px}.fact{border:1px solid var(--line-dim);padding:4px 6px;background:var(--panel-soft)}.fact b{color:var(--muted);font-weight:700}.history{min-height:240px}.range{display:flex;gap:5px;align-items:center;flex-wrap:wrap}.pill{border:1px solid var(--line);padding:2px 7px;border-radius:999px;background:var(--panel-soft);color:var(--text)}.pill.active{border-color:var(--accent);color:var(--accent);font-weight:900}.bars{display:grid;grid-template-columns:repeat(14,1fr);gap:3px;align-items:end;height:82px;margin:8px 0 2px;border-bottom:1px solid var(--line)}.bar{display:flex;flex-direction:column;justify-content:flex-end;align-items:center;min-width:0}.bar span{display:block;width:100%;text-align:center}.allocated{color:var(--graph)}.unalloc{color:var(--unalloc)}.gap{color:var(--gap)}.zero{color:var(--zero)}.partial{border-top:2px dashed var(--select)}.selected{outline:1px solid var(--select);outline-offset:1px}.day-row{display:grid;grid-template-columns:repeat(14,1fr);gap:3px;color:var(--muted);font-size:12px;text-align:center}.marker{color:var(--select);text-align:center;letter-spacing:.5ch}.legend{color:var(--muted);font-size:12px;margin-top:7px}.readout{border-top:1px dashed var(--line);padding-top:7px;margin-top:7px}.health{display:grid;grid-template-rows:auto auto;gap:8px}.subbox{border:1px solid var(--line-dim);padding:7px 8px;background:var(--panel-soft)}.subbox h3{margin:0 0 5px;font-size:14px;color:var(--accent)}.settings-grid{display:grid;grid-template-columns:auto 1fr;gap:3px 9px}.setting-control{color:var(--accent);font-weight:900}.service{margin-top:8px;border:1px solid var(--line);background:var(--panel);display:grid;grid-template-columns:1fr auto;gap:8px;padding:8px 10px}.continuity{margin-top:3px}.actions{display:flex;align-items:center;gap:10px;white-space:nowrap}.quit{border:2px solid var(--text);padding:3px 8px;font-weight:900;color:var(--text);background:var(--quit-bg)}.long-reason{display:none;color:var(--bad);margin-top:5px}body[data-state=error] .long-reason,body[data-state=stale] .long-reason,body[data-state=interrupted] .long-reason{display:block}.narrow-only{display:none}.switcher{position:fixed;z-index:10;bottom:14px;left:50%;transform:translateX(-50%);display:flex;gap:8px;align-items:center;flex-wrap:wrap;justify-content:center;max-width:calc(100vw - 20px);background:#f5f7fb;color:#151923;border-radius:999px;padding:8px 11px;box-shadow:0 8px 32px #000}.switcher label{font-weight:800}.switcher select,.switcher button{font:700 13px ui-monospace,monospace;border:1px solid #c8ced9;border-radius:999px;background:white;color:#151923;padding:6px 9px}.switcher button.active{background:#151923;color:white}.switcher .check{display:flex;align-items:center;gap:4px}.switcher input{accent-color:#151923}
|
||||
body[data-preset=amber]{--bg:#0d0b07;--panel:#17120b;--panel-soft:#20180f;--title-bg:#22170a;--title:#f8e1ad;--text:#f3eadb;--muted:#c2ad8a;--line:#8a6432;--line-dim:#5a4022;--accent:#f0b35a;--graph:#f0b35a;--unalloc:#d38f3f;--gap:#6f604e;--zero:#d8c5a3;--select:#f8e08e;--ok:#5ee08b;--warn:#ffd166;--bad:#ff6b6b;--neutral:#b8bec9;--warn-bg:#3a2c0c;--warn-text:#fff3c7;--quit-bg:#2a1e10;--term-w:960px;--term-h:650px}
|
||||
body[data-preset=nord]{--bg:#0b1118;--panel:#111827;--panel-soft:#162033;--title-bg:#172235;--title:#e5e9f0;--text:#e5e9f0;--muted:#a9b4c8;--line:#4c566a;--line-dim:#354154;--accent:#88c0d0;--graph:#88c0d0;--unalloc:#d08770;--gap:#677285;--zero:#c5ccd8;--select:#ebcb8b;--ok:#5ee08b;--warn:#ffd166;--bad:#ff6b6b;--neutral:#b8bec9;--warn-bg:#2f2a17;--warn-text:#fff0bd;--quit-bg:#172235;--term-w:960px;--term-h:650px}
|
||||
body[data-preset=high]{--bg:#000;--panel:#000;--panel-soft:#111;--title-bg:#000;--title:#fff;--text:#fff;--muted:#d8d8d8;--line:#fff;--line-dim:#888;--accent:#00ffff;--graph:#ffff00;--unalloc:#ffaf00;--gap:#9b9b9b;--zero:#fff;--select:#00ffff;--ok:#00ff00;--warn:#ffff00;--bad:#ff4040;--neutral:#fff;--warn-bg:#333300;--warn-text:#fff;--quit-bg:#000;--term-w:960px;--term-h:650px}
|
||||
body[data-width=wide]{--term-w:1180px;--term-h:760px}.wide-note{display:none}body[data-width=wide] .wide-note{display:inline}
|
||||
body[data-width=narrow]{--term-w:620px;--term-h:600px;font-size:13px}body[data-width=narrow] .grid{grid-template-columns:1fr}body[data-width=narrow] .headline{grid-column:auto}body[data-width=narrow] .facts{grid-template-columns:repeat(2,minmax(0,1fr))}body[data-width=narrow] .bars,body[data-width=narrow] .day-row,body[data-width=narrow] .marker,body[data-width=narrow] .legend{display:none}body[data-width=narrow] .narrow-only{display:block;border:1px dashed var(--line);padding:8px;margin:8px 0;color:var(--warn)}body[data-width=narrow] .history{min-height:auto}body[data-width=narrow] .service{grid-template-columns:1fr}body[data-width=narrow] .actions{white-space:normal;justify-content:space-between}.sr{position:absolute;left:-10000px}
|
||||
</style>
|
||||
</head>
|
||||
<body data-preset="amber" data-state="monitoring" data-width="normal" data-glyph="emoji">
|
||||
<p class="note"><b>PROTOTYPE:</b> titlebox, health/settings grouping, theme presets and status accessibility for <i>Approve titlebox, health layout and colour presets</i>. Static throwaway; no production code.</p>
|
||||
<main class="terminal" aria-label="Fenris terminal mockup">
|
||||
<div class="screen">
|
||||
<div class="titlebox" aria-label="Application title"><span class="emoji">🐺 Fenris by Bongbetic</span><span class="fallback">Fenris by Bongbetic</span></div>
|
||||
<div class="auth">privileged actions will prompt for authentication (polkit) <span class="wide-note">· clears after first refresh</span></div>
|
||||
<section class="pane headline">
|
||||
<div class="pane-title">lifespan + monitoring</div>
|
||||
<div class="big">Usage-adjusted theoretical lifespan: <span class="accent">8 yr 103 d remaining</span></div>
|
||||
<div>if current habits continue · sustained regime: 46 days · scenario range: 7d 12.4 / 28d 11.1 / 90d 10.7 GB/day</div>
|
||||
<div class="status-line"><div class="state monitoring" id="stateLine"><span class="glyph">●</span><span id="stateLabel">Monitoring</span></div><div id="stateReason">Last sample: 4 min ago · blink means enabled + fresh, not data refresh</div></div>
|
||||
<div class="long-reason" id="longReason">last run failed (exit 3) · last good sample 4 min ago · collection stopped outside Fenris details stay visible without changing the status vocabulary</div>
|
||||
<div class="pause-banner">‖ Paused — monitoring paused · paused time excluded from your usage habit</div>
|
||||
<div class="facts"><div class="fact"><b>Boot</b><br><span id="bootFact">on</span></div><div class="fact"><b>Runtime</b><br><span id="runtimeFact">timer active</span></div><div class="fact"><b>Last collect</b><br><span id="collectFact">ok · 4 min ago</span></div><div class="fact"><b>Freshness</b><br><span id="freshFact">fresh</span></div></div>
|
||||
</section>
|
||||
<div class="grid">
|
||||
<section class="pane history" tabindex="0" aria-label="Usage history graph pane">
|
||||
<div class="pane-title">usage history</div>
|
||||
<div class="range"><b>usage history · Local · UTC+05:30 · Asia/Kolkata</b><span class="pill">7</span><span class="pill active">14</span><span class="pill">28</span><span class="pill">90</span></div>
|
||||
<div class="narrow-only">14-day write history: 9.2–18.6 GB/day · selected Wed 09 Sep · graph needs ≥80×24</div>
|
||||
<div class="bars" aria-hidden="true">
|
||||
<div class="bar"><span class="gap">░</span><span class="gap">░</span></div><div class="bar"><span class="zero">·</span></div><div class="bar"><span class="allocated">█</span><span class="allocated">█</span></div><div class="bar"><span class="allocated">█</span><span class="unalloc">▒</span><span class="unalloc">▒</span></div><div class="bar selected"><span class="allocated">█</span><span class="allocated">█</span><span class="allocated">█</span><span class="unalloc">▒</span></div><div class="bar"><span class="allocated">█</span></div><div class="bar"><span class="allocated">█</span><span class="allocated">█</span></div><div class="bar"><span class="gap">░</span><span class="gap">░</span><span class="gap">░</span></div><div class="bar"><span class="allocated">█</span><span class="allocated">█</span><span class="allocated">█</span><span class="allocated">█</span></div><div class="bar"><span class="allocated">█</span><span class="unalloc">▒</span></div><div class="bar"><span class="zero">·</span></div><div class="bar"><span class="allocated">█</span><span class="allocated">█</span></div><div class="bar partial"><span class="allocated">█</span><span class="allocated">█</span><span class="unalloc">▒</span></div><div class="bar partial"><span class="allocated">█</span></div>
|
||||
</div>
|
||||
<div class="marker"> ▼</div>
|
||||
<div class="day-row"><span>27</span><span>28</span><span>29</span><span>30</span><span>31</span><span>01</span><span>02</span><span>03</span><span>04</span><span>05</span><span>06</span><span>07</span><span>08</span><span>09</span></div>
|
||||
<div class="legend">legend: <span class="allocated">█ allocated</span> · <span class="unalloc">▒ unallocated</span> · <span class="gap">░ gap/no evidence</span> · <span class="zero">· 0 B</span> · ┄ partial · ▼ selected</div>
|
||||
<div class="readout"><b>Wed 09 Sep</b> · writes 14.2 GB allocated + 2.0 GB unallocated · 21 evidenced hours · coverage 87% · partial · 13 h elapsed · so far<br><span class="muted">Enter/click drills into hourly detail; Esc/Backspace returns. Footer shows graph keys while graph is focused.</span></div>
|
||||
</section>
|
||||
<section class="pane health" tabindex="0" aria-label="Drive health and settings pane">
|
||||
<div class="subbox"><h3>Drive health</h3><b>Samsung SSD 990 PRO</b><br>temperature 41°C · spare 100%<br>media errors 0 · unsafe shutdowns 2<br>power-on 1,184 h · 46 cycles · 2 TB<br><b>vendor wear</b>: 3% used · 18.4 TB written <span class="muted">(context, not a second projection)</span></div>
|
||||
<div class="subbox"><h3>Settings</h3><div class="settings-grid"><span>device</span><span>/dev/disk/by-id/nvme-Samsung_990</span><span>baseline</span><span>verified override · 1,200 TBW</span><span>retention</span><span>raw 14 d · intervals/hours/days indefinite</span><span>theme</span><span class="setting-control" id="themeText">Amber preset · user preference</span><span>motion</span><span class="setting-control" id="motionText">normal blink</span></div></div>
|
||||
</section>
|
||||
</div>
|
||||
<section class="service" aria-label="Service strip and actions">
|
||||
<div><b>CONTINUITY</b> <span id="continuity">monitoring: active in background · persists across reboots</span><div class="continuity muted">p pause · r resume · c collect · d disclosures · graph focus: ←/→ select · Enter drill · 1/2/3/4 ranges · t preset · m motion</div></div>
|
||||
<div class="actions"><span class="quit">q QUIT TUI</span></div>
|
||||
</section>
|
||||
</div>
|
||||
</main>
|
||||
<nav class="switcher" aria-label="Prototype controls">
|
||||
<label>preset <select id="preset"><option value="amber">Amber</option><option value="nord">Nord</option><option value="high">High Contrast</option></select></label>
|
||||
<label>state <select id="state"><option value="monitoring">Monitoring</option><option value="collecting">Collecting</option><option value="paused">Paused</option><option value="waiting">Waiting</option><option value="interrupted">Interrupted</option><option value="error">Error</option><option value="stale">Stale</option><option value="unknown">Unknown</option></select></label>
|
||||
<label>size <select id="width"><option value="normal">80×24</option><option value="wide">140×40</option><option value="narrow">narrow</option></select></label>
|
||||
<button id="glyph">wolf title</button>
|
||||
<label class="check"><input id="motion" type="checkbox"> reduced motion</label>
|
||||
</nav>
|
||||
<script>
|
||||
const states={
|
||||
monitoring:{glyph:'●',label:'Monitoring',cls:'monitoring',reason:'Last sample: 4 min ago · blink means enabled + fresh, not data refresh',boot:'on',runtime:'timer active',collect:'ok · 4 min ago',fresh:'fresh',cont:'monitoring: active in background · persists across reboots'},
|
||||
collecting:{glyph:'◐',label:'Collecting',cls:'collecting',reason:'run in flight (≤90 s) · underlying state: monitoring',boot:'on',runtime:'service activating',collect:'running now',fresh:'fresh',cont:'monitoring: active in background · persists across reboots'},
|
||||
paused:{glyph:'‖',label:'Paused',cls:'paused',reason:'monitoring paused — paused time excluded from your usage habit',boot:'off',runtime:'timer inactive',collect:'ok before pause',fresh:'last sample 49 days ago',cont:'monitoring: does not start on next boot'},
|
||||
waiting:{glyph:'○',label:'Waiting',cls:'waiting',reason:'awaiting another sample · first graph data is independent of lifespan warm-up',boot:'on',runtime:'timer active',collect:'ok · one sample',fresh:'empty / insufficient pair',cont:'monitoring: active in background · persists across reboots'},
|
||||
interrupted:{glyph:'⊘',label:'Interrupted',cls:'interrupted',reason:'collection stopped outside Fenris — monitoring period still open',boot:'on',runtime:'timer inactive',collect:'unknown stop',fresh:'last sample 3 days ago',cont:'enabled but not currently collecting'},
|
||||
error:{glyph:'✖',label:'Error',cls:'error',reason:'last run failed (exit 3) · last good sample 4 min ago',boot:'on',runtime:'timer active',collect:'FAILED exit 3',fresh:'fresh data, failed last run',cont:'monitoring: active in background · persists across reboots'},
|
||||
stale:{glyph:'◌',label:'Stale',cls:'stale',reason:'last sample 3 days ago · timer active, no failure recorded',boot:'on',runtime:'timer active',collect:'ok then silent',fresh:'stale',cont:'monitoring: active in background · persists across reboots'},
|
||||
unknown:{glyph:'?',label:'Unknown',cls:'unknown',reason:'service state unavailable',boot:'unknown',runtime:'unknown',collect:'unknown',fresh:'store facts unavailable',cont:'service state unavailable'}
|
||||
};
|
||||
function apply(){
|
||||
const preset=document.getElementById('preset').value; const state=document.getElementById('state').value; const width=document.getElementById('width').value;
|
||||
document.body.dataset.preset=preset; document.body.dataset.state=state; document.body.dataset.width=width;
|
||||
const s=states[state]; const line=document.getElementById('stateLine'); line.className='state '+s.cls; line.querySelector('.glyph').textContent=s.glyph; document.getElementById('stateLabel').textContent=s.label; document.getElementById('stateReason').textContent=s.reason;
|
||||
document.getElementById('bootFact').textContent=s.boot; document.getElementById('runtimeFact').textContent=s.runtime; document.getElementById('collectFact').textContent=s.collect; document.getElementById('freshFact').textContent=s.fresh; document.getElementById('continuity').textContent=s.cont;
|
||||
document.getElementById('themeText').textContent={amber:'Amber preset · user preference',nord:'Nord preset · user preference',high:'High Contrast preset · user preference'}[preset];
|
||||
const reduced=document.getElementById('motion').checked; document.body.classList.toggle('reduced',reduced); document.getElementById('motionText').textContent=reduced?'reduced motion · steady status dot':'normal blink';
|
||||
const params=new URLSearchParams({preset,state,width,glyph:document.body.dataset.glyph,motion:reduced?'reduced':'normal'}); history.replaceState(null,'','?'+params);
|
||||
}
|
||||
const params=new URLSearchParams(location.search); for(const id of ['preset','state','width']) if(params.get(id)) document.getElementById(id).value=params.get(id)); if(params.get('glyph')==='plain') document.body.dataset.glyph='plain'; if(params.get('motion')==='reduced') document.getElementById('motion').checked=true;
|
||||
document.getElementById('glyph').onclick=()=>{document.body.dataset.glyph=document.body.dataset.glyph==='emoji'?'plain':'emoji'; document.getElementById('glyph').textContent=document.body.dataset.glyph==='emoji'?'wolf title':'plain title'; apply()};
|
||||
for(const id of ['preset','state','width','motion']) document.getElementById(id).addEventListener('change',apply);
|
||||
apply();
|
||||
</script>
|
||||
</body>
|
||||
</html>
|
||||
@@ -1,6 +0,0 @@
|
||||
#!/bin/sh
|
||||
# PROTOTYPE: static local server; no dependencies or persistence.
|
||||
set -eu
|
||||
cd "$(git rev-parse --show-toplevel)"
|
||||
printf '%s\n' 'Open http://127.0.0.1:8765/prototype/titlebox-health-presets/?preset=amber&state=monitoring'
|
||||
exec python3 -m http.server 8765 --bind 127.0.0.1
|
||||
+40
-5
@@ -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,
|
||||
|
||||
@@ -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()
|
||||
+133
-2
@@ -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
|
||||
|
||||
@@ -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
|
||||
+24
-2
@@ -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)
|
||||
@@ -598,6 +609,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 +659,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
@@ -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}")
|
||||
|
||||
+24
-7
@@ -447,15 +447,32 @@ 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
|
||||
|
||||
# --- 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) ---
|
||||
history = _query_usage_history(conn)
|
||||
|
||||
@@ -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",
|
||||
}
|
||||
+54
-3
@@ -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
|
||||
|
||||
@@ -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"
|
||||
@@ -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
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
Reference in New Issue
Block a user