242 lines
8.6 KiB
Python
242 lines
8.6 KiB
Python
"""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, local_tz "
|
|
"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, local_tz "
|
|
"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], "local_tz": row[8],
|
|
}
|
|
|
|
|
|
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.
|
|
Caller owns the transaction.
|
|
"""
|
|
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, bytes_read_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_br = existing[2] + br_delta
|
|
new_samples = existing[3] + sample_count
|
|
conn.execute(
|
|
"UPDATE hour_observations "
|
|
"SET bytes_written_delta = ?, bytes_read_delta = ?, sample_count = ? "
|
|
"WHERE id = ?",
|
|
(new_bw, new_br, new_samples, existing[0]),
|
|
)
|
|
def _add_unattributed_bytes(
|
|
conn: sqlite3.Connection,
|
|
prev_ts: datetime,
|
|
next_ts: datetime,
|
|
bw_delta: int,
|
|
br_delta: int,
|
|
) -> None:
|
|
"""Store unattributed byte deltas once as shared boundary evidence.
|
|
|
|
A cross-midnight interval's delta is preserved on the day where it
|
|
STARTS (the earlier day). It is not duplicated into both days;
|
|
the spec requires preserving the measured volume once as shared
|
|
unallocated boundary evidence (issue #88).
|
|
"""
|
|
day = prev_ts.strftime("%Y-%m-%d")
|
|
|
|
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),
|
|
)
|