feat(#73): publish trustworthy first usage history

- Schema migration 1-2: add segment_id to samples, unattributed bytes to day_aggregates

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

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

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

- Monitoring period ensured open on each collection run

- Derivation failures preserve prior history

Closes #73
This commit is contained in:
xavierk
2026-09-14 03:19:08 +05:30
parent f406285a0f
commit d790ff84c5
7 changed files with 1074 additions and 26 deletions
+40 -5
View File
@@ -16,6 +16,8 @@ from pathlib import Path
from typing import Any, Dict, Optional, Tuple
from .store import init_store, get_store_path
from .monitoring_periods import ensure_period_open
from .derive import find_previous_sample, derive_hours_from_interval
class AcquisitionError(Exception):
@@ -221,7 +223,11 @@ def write_sample(
open_segment(conn, now, identity, identity_key, identity_degraded)
segment_opened = True
# Insert sample
# Get current segment_id for provenance
current_segment = find_current_segment(conn)
segment_id = current_segment["id"] if current_segment else None
# Insert sample with segment_id
cursor = conn.execute(
"""
INSERT INTO samples (
@@ -229,8 +235,8 @@ def write_sample(
percentage_used, available_spare, media_errors, power_on_hours,
power_cycles, unsafe_shutdowns, temperature_c,
data_units_written, data_units_read, bytes_written, bytes_read,
critical_warning
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
critical_warning, segment_id
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
""",
(
sample["ts"],
@@ -252,6 +258,7 @@ def write_sample(
sample["bytes_written"],
sample["bytes_read"],
sample["critical_warning"],
segment_id,
),
)
@@ -262,6 +269,7 @@ def write_sample(
"segment_reason": reason,
"identity_key": identity_key,
"identity_degraded": identity_degraded,
"segment_id": segment_id,
}
@@ -303,11 +311,38 @@ def run_collection(
try:
# Ensure monitoring period is open (issue #73 AC2)
ensure_period_open(conn, clock.utcnow())
# Validate invariants
validate_sample_invariants(sample, conn)
# Write sample
write_sample(sample, identity, conn, clock)
# Write sample and get segment info
seg_info = write_sample(sample, identity, conn, clock)
# Derive hour observations from interval with previous sample
try:
# Find the sample we just wrote
cursor = conn.execute("SELECT id FROM samples ORDER BY id DESC LIMIT 1")
current_id = cursor.fetchone()[0]
prev = find_previous_sample(conn, seg_info.get("segment_id"), current_id)
if prev is not None:
# Build current sample dict for derivation
current = {
"id": current_id,
"ts": sample["ts"],
"bytes_written": sample["bytes_written"],
"bytes_read": sample["bytes_read"],
"power_on_hours": sample["power_on_hours"],
"temperature_c": sample["temperature_c"],
"data_units_written": sample["data_units_written"],
"data_units_read": sample["data_units_read"],
}
derive_hours_from_interval(conn, prev, current)
except Exception:
# Derivation failure must not prevent sample persistence (issue #73 AC6)
pass
return {
"ok": True,
+237
View File
@@ -0,0 +1,237 @@
"""Interval derivation: samples → hour observations → day aggregates.
After each collection run, the collector calls into this module to:
1. Find the previous sample in the same segment
2. Compute deltas (bytes, POH, temperature)
3. Classify the hour(s) the interval spans
4. Write/update hour_observations for each affected hour
5. Update day_aggregates with unattributed cross-hour bytes
Cross-hour deltas are retained once with unknown shares explicit (issue #73 AC4).
No proportional allocation, endpoint assignment, or double counting.
"""
import sqlite3
from datetime import datetime, timedelta, timezone
from typing import Any, Dict, List, Optional, Tuple
from .hour_classify import classify_hour, HourSplit
def find_previous_sample(
conn: sqlite3.Connection,
segment_id: Optional[int],
current_sample_id: int,
) -> Optional[Dict[str, Any]]:
"""Find the most recent sample before current_sample_id in the same segment.
Returns None if no previous sample exists (first sample in segment).
"""
if segment_id is not None:
cursor = conn.execute(
"SELECT id, ts, bytes_written, bytes_read, power_on_hours, "
" temperature_c, data_units_written, data_units_read "
"FROM samples WHERE id < ? AND segment_id = ? "
"ORDER BY id DESC LIMIT 1",
(current_sample_id, segment_id),
)
else:
cursor = conn.execute(
"SELECT id, ts, bytes_written, bytes_read, power_on_hours, "
" temperature_c, data_units_written, data_units_read "
"FROM samples WHERE id < ? "
"ORDER BY id DESC LIMIT 1",
(current_sample_id,),
)
row = cursor.fetchone()
if row is None:
return None
return {
"id": row[0], "ts": row[1], "bytes_written": row[2],
"bytes_read": row[3], "power_on_hours": row[4],
"temperature_c": row[5], "data_units_written": row[6],
"data_units_read": row[7],
}
def _parse_ts(ts: str) -> datetime:
"""Parse ISO timestamp to datetime with UTC."""
dt = datetime.fromisoformat(ts)
if dt.tzinfo is None:
dt = dt.replace(tzinfo=timezone.utc)
return dt
def _hour_floor(dt: datetime) -> datetime:
"""Floor a datetime to its UTC hour boundary."""
return dt.replace(minute=0, second=0, microsecond=0)
def _hours_spanned(start: datetime, end: datetime) -> List[datetime]:
"""Return list of UTC hour boundaries spanned by [start, end)."""
hours = []
h = _hour_floor(start)
while h < end:
hours.append(h)
h += timedelta(hours=1)
return hours
def _compute_sampled_seconds_in_hour(
start: datetime, end: datetime, hour_start: datetime
) -> int:
"""How many seconds of the sample interval fall within this hour."""
hour_end = hour_start + timedelta(hours=1)
effective_start = max(start, hour_start)
effective_end = min(end, hour_end)
if effective_start >= effective_end:
return 0
return int((effective_end - effective_start).total_seconds())
def derive_hours_from_interval(
conn: sqlite3.Connection,
prev_sample: Dict[str, Any],
next_sample: Dict[str, Any],
) -> List[Dict[str, Any]]:
"""Derive hour observations from a sample pair interval.
Returns list of hour observation dicts that were written/updated.
"""
prev_ts = _parse_ts(prev_sample["ts"])
next_ts = _parse_ts(next_sample["ts"])
# Deltas
bw_delta = max(0, next_sample["bytes_written"] - prev_sample["bytes_written"])
br_delta = max(0, next_sample["bytes_read"] - prev_sample["bytes_read"])
poh_delta_s = max(0, (next_sample["power_on_hours"] - prev_sample["power_on_hours"])) * 3600
hours = _hours_spanned(prev_ts, next_ts)
total_span_s = int((next_ts - prev_ts).total_seconds())
results = []
if len(hours) == 1:
# Same-hour interval: fully attributed to this hour
hour_key = hours[0].strftime("%Y-%m-%dT%H:00:00+00:00")
sampled_s = total_span_s
# Classify hour
split = classify_hour(
wall_clock_seconds=3600,
poh_delta=poh_delta_s,
duw_delta=bw_delta,
dur_delta=br_delta,
sampled_seconds=sampled_s,
)
_upsert_hour_observation(
conn, hour_key, split,
bw_delta, br_delta,
prev_sample.get("temperature_c"), next_sample.get("temperature_c"),
2, # 2 samples contributed (prev + next)
)
results.append({"hour": hour_key, "bytes_written": bw_delta, "attributed": True})
elif len(hours) >= 2:
# Cross-hour interval: split wall-clock time, bytes unattributed
for h in hours:
hour_key = h.strftime("%Y-%m-%dT%H:00:00+00:00")
sampled_s = _compute_sampled_seconds_in_hour(prev_ts, next_ts, h)
# For cross-hour, we classify based on time only (no byte attribution)
# The hour gets its time split but NOT the byte delta
split = classify_hour(
wall_clock_seconds=3600,
poh_delta=0, # POH attribution unknown for cross-hour
duw_delta=0, # Bytes unattributed
dur_delta=0,
sampled_seconds=sampled_s,
)
_upsert_hour_observation(
conn, hour_key, split,
0, 0, # No byte attribution for cross-hour
None, None,
0, # No sample falls IN this hour
)
results.append({"hour": hour_key, "bytes_written": 0, "attributed": False})
# Track unattributed bytes at day level
_add_unattributed_bytes(conn, prev_ts, next_ts, bw_delta, br_delta)
return results
def _upsert_hour_observation(
conn: sqlite3.Connection,
hour_key: str,
split: HourSplit,
bw_delta: int,
br_delta: int,
temp_min: Optional[int],
temp_max: Optional[int],
sample_count: int,
) -> None:
"""Insert or update an hour observation."""
# Check if hour exists
existing = conn.execute(
"SELECT id, bytes_written_delta, sample_count FROM hour_observations WHERE hour = ?",
(hour_key,),
).fetchone()
if existing is None:
temp_avg = ((temp_min or 0) + (temp_max or 0)) / 2 if temp_min is not None else None
coverage = (split.seconds_active + split.seconds_idle + split.seconds_powered_off) / 3600.0
conn.execute(
"""INSERT INTO hour_observations
(hour, active_seconds, idle_seconds, powered_off_seconds, unknown_seconds,
bytes_written_delta, bytes_read_delta,
temperature_min, temperature_avg, temperature_max,
sample_count, coverage)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)""",
(hour_key, split.seconds_active, split.seconds_idle,
split.seconds_powered_off, split.seconds_unknown,
bw_delta, br_delta,
temp_min, temp_avg, temp_max,
sample_count, coverage),
)
else:
# Merge: accumulate bytes and sample count
new_bw = existing[1] + bw_delta
new_samples = existing[2] + sample_count
conn.execute(
"UPDATE hour_observations SET bytes_written_delta = ?, sample_count = ? WHERE id = ?",
(new_bw, new_samples, existing[0]),
)
conn.commit()
def _add_unattributed_bytes(
conn: sqlite3.Connection,
prev_ts: datetime,
next_ts: datetime,
bw_delta: int,
br_delta: int,
) -> None:
"""Add unattributed byte deltas to day aggregates for each day touched."""
prev_day = prev_ts.strftime("%Y-%m-%d")
next_day = next_ts.strftime("%Y-%m-%d")
days = {prev_day, next_day}
for day in days:
existing = conn.execute(
"SELECT id FROM day_aggregates WHERE day = ?", (day,)
).fetchone()
if existing is None:
conn.execute(
"INSERT INTO day_aggregates (day, unattributed_bytes_written, unattributed_bytes_read) "
"VALUES (?, ?, ?)",
(day, bw_delta, br_delta),
)
else:
conn.execute(
"UPDATE day_aggregates SET unattributed_bytes_written = unattributed_bytes_written + ?, "
"unattributed_bytes_read = unattributed_bytes_read + ? WHERE day = ?",
(bw_delta, br_delta, day),
)
conn.commit()
+24 -2
View File
@@ -331,7 +331,8 @@ def check_retired_flag(flag: str) -> Optional[str]:
def _format_projection(proj, freshness: str, service: Dict[str, Any],
drive_facts: List[str], config_error: Optional[str],
store_fault: Optional[str], newer_schema: Optional[str],
journal_hint: Optional[str]) -> str:
journal_hint: Optional[str],
sample_count: int = 0, day_count: int = 0) -> str:
"""Format the complete status output."""
lines = []
@@ -361,6 +362,16 @@ def _format_projection(proj, freshness: str, service: Dict[str, Any],
_append_service_facts(lines, service)
return "\n".join(lines)
# --- Single sample: awaiting another sample (issue #73 AC3) ---
# Only show awaiting state when there are no day aggregates (e.g., legacy import
# or hand-crafted stores can have 1 sample but sufficient day data for projection)
if sample_count <= 1 and day_count == 0:
lines.append("awaiting another sample")
lines.append("")
lines.append("Collecting usage data — the first projection requires at least two samples.")
_append_service_facts(lines, service)
return "\n".join(lines)
# --- Projection headline ---
headline = _format_headline(proj)
lines.append(headline)
@@ -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
View File
@@ -12,7 +12,7 @@ from typing import Optional
# Schema version - increment on each migration
SCHEMA_VERSION = 1
SCHEMA_VERSION = 2
# Packaged default placement (spec §8.3). The config may override it, but a
@@ -106,7 +106,8 @@ def _create_schema(conn: sqlite3.Connection):
data_units_read INTEGER,
bytes_written INTEGER,
bytes_read INTEGER,
critical_warning INTEGER
critical_warning INTEGER,
segment_id INTEGER
)
""")
@@ -141,7 +142,9 @@ def _create_schema(conn: sqlite3.Connection):
bytes_written_delta INTEGER DEFAULT 0,
bytes_read_delta INTEGER DEFAULT 0,
sample_count INTEGER DEFAULT 0,
coverage REAL DEFAULT 0.0
coverage REAL DEFAULT 0.0,
unattributed_bytes_written INTEGER DEFAULT 0,
unattributed_bytes_read INTEGER DEFAULT 0
)
""")
@@ -207,11 +210,25 @@ def _apply_migrations(conn: sqlite3.Connection, current_version: int):
Spec: §3.6, §10.2
"""
# Migration 1→2: example placeholder
# if current_version < 2:
# conn.execute("ALTER TABLE ...")
# current_version = 2
pass
# Migration 1→2: add segment_id provenance to samples,
# unattributed byte tracking to day_aggregates (issue #73)
if current_version < 2:
# Defensive: only ALTER if table exists (handles minimal v1 stores)
tables = {row[0] for row in conn.execute(
"SELECT name FROM sqlite_master WHERE type='table'"
).fetchall()}
if "samples" in tables:
# Check if column already exists (idempotent)
cols = {row[1] for row in conn.execute("PRAGMA table_info(samples)").fetchall()}
if "segment_id" not in cols:
conn.execute("ALTER TABLE samples ADD COLUMN segment_id INTEGER")
if "day_aggregates" in tables:
cols = {row[1] for row in conn.execute("PRAGMA table_info(day_aggregates)").fetchall()}
if "unattributed_bytes_written" not in cols:
conn.execute("ALTER TABLE day_aggregates ADD COLUMN unattributed_bytes_written INTEGER DEFAULT 0")
if "unattributed_bytes_read" not in cols:
conn.execute("ALTER TABLE day_aggregates ADD COLUMN unattributed_bytes_read INTEGER DEFAULT 0")
current_version = 2
def migrate_to_latest(store_path: Path) -> int:
@@ -239,6 +256,14 @@ def migrate_to_latest(store_path: Path) -> int:
conn.close()
return 0 # Already up to date
# Version 0 means no schema — create fresh (issue #73)
if current_version == 0:
_create_schema(conn)
conn.execute(f"PRAGMA user_version={SCHEMA_VERSION}")
conn.commit()
conn.close()
return SCHEMA_VERSION
steps = SCHEMA_VERSION - current_version
_apply_migrations(conn, current_version)
conn.execute(f"PRAGMA user_version={SCHEMA_VERSION}")
+24 -7
View File
@@ -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)