diff --git a/AGENTS.md b/AGENTS.md index 6166ac3c0..ad276acf1 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -203,6 +203,7 @@ Each domain has exactly **one** write-owning module (or one tightly-scoped famil | Schedules (`config/schedules.json`) | `solstone/think/schedule_config.py` | | Push devices (`config/push_devices.json`) | `solstone/think/push/devices.py` | | Local inference operational telemetry (`health/local-inference/YYYYMMDD.jsonl`) | `solstone/think/providers/local_admission.py` | +| Media offload ledger (`health/offload/.jsonl`) | `solstone/think/offload_ledger.py` | | Parakeet server placement record (`health/parakeet-cpp.placement`) | `solstone/think/providers/parakeet_server.py` | | Hosted backup binding (`backup/hosted/binding.json`) | `solstone/think/backup/hosted.py` | | Convey config (`config/convey.json`) | `solstone/convey/config.py` + `solstone/think/facets.py` | diff --git a/scripts/check_journal_io_access.py b/scripts/check_journal_io_access.py index fb1715b58..dbcb19914 100644 --- a/scripts/check_journal_io_access.py +++ b/scripts/check_journal_io_access.py @@ -125,6 +125,7 @@ OWNER_FILES: frozenset[str] = frozenset( "solstone/think/identity.py", "solstone/think/journal_config.py", "solstone/think/log_retention.py", + "solstone/think/offload_ledger.py", # Sole writer of content-free bundled-local inference telemetry. "solstone/think/providers/local_admission.py", "solstone/think/schedule_config.py", diff --git a/solstone/think/backup/state.py b/solstone/think/backup/state.py index b781f389f..ba7e9f129 100644 --- a/solstone/think/backup/state.py +++ b/solstone/think/backup/state.py @@ -39,6 +39,11 @@ BACKUP_DEFAULTS: dict[str, Any] = { "weekly": 4, "monthly": 12, }, + "offload": { + "enabled": False, + "budget_bytes": None, + "floor_bytes": None, + }, "schedule": { "every": "daily", "enabled": False, @@ -54,8 +59,23 @@ BACKUP_DEFAULTS: dict[str, Any] = { "status": None, "error_reason": None, }, + # Offload records use "reason" instead of "error_reason" because skipped + # and stalled are expected outcomes, not errors. + "last_offload": { + "time": None, + "status": None, + "reason": None, + }, + "last_verification": { + "time": None, + "status": None, + "reason": None, + }, } +OFFLOAD_KEYS = ("enabled", "budget_bytes", "floor_bytes") +OFFLOAD_STATUSES = ("ok", "skipped", "stalled", "error") RETENTION_KEYS = ("hourly", "daily", "weekly", "monthly") +VERIFICATION_STATUSES = ("ok", "skipped", "error") @dataclass(frozen=True) @@ -200,6 +220,33 @@ def set_retention(retention: dict[str, int]) -> None: write_journal_config(config) +def set_offload(offload: dict[str, Any]) -> None: + if not isinstance(offload, dict): + raise ValueError("backup offload must be a JSON object") + if set(offload) != set(OFFLOAD_KEYS): + raise ValueError( + "backup offload must include enabled, budget_bytes, floor_bytes" + ) + if not isinstance(offload["enabled"], bool): + raise ValueError("backup offload enabled must be a boolean") + for key in ("budget_bytes", "floor_bytes"): + value = offload[key] + if value is not None and (type(value) is not int or value <= 0): + raise ValueError( + "backup offload byte values must be positive integers or null" + ) + + with hold_config_lock(): + config = read_journal_config() + backup = _writable_backup_section(config) + backup["offload"] = { + "enabled": offload["enabled"], + "budget_bytes": offload["budget_bytes"], + "floor_bytes": offload["floor_bytes"], + } + write_journal_config(config) + + def set_recovery_key(recovery_key: str) -> None: with hold_config_lock(): config = read_journal_config() @@ -251,6 +298,56 @@ def record_prune_result( write_journal_config(config) +def record_offload_result( + *, + status: str, + time: int | None, + reason: str | None = None, +) -> None: + """Record the last media-offload run. + + Convention: reason is None on ok. This layer only enforces the closed + status vocabulary; callers own richer run-state validation. + """ + if status not in OFFLOAD_STATUSES: + raise ValueError("backup offload status must be ok, skipped, stalled, or error") + + with hold_config_lock(): + config = read_journal_config() + backup = _writable_backup_section(config) + backup["last_offload"] = { + "time": time, + "status": status, + "reason": reason, + } + write_journal_config(config) + + +def record_verification_result( + *, + status: str, + time: int | None, + reason: str | None = None, +) -> None: + """Record the last media-offload verification run. + + Convention: reason is None on ok. This layer only enforces the closed + status vocabulary; callers own richer run-state validation. + """ + if status not in VERIFICATION_STATUSES: + raise ValueError("backup verification status must be ok, skipped, or error") + + with hold_config_lock(): + config = read_journal_config() + backup = _writable_backup_section(config) + backup["last_verification"] = { + "time": time, + "status": status, + "reason": reason, + } + write_journal_config(config) + + def status_view() -> dict[str, Any]: config = get_backup_config() destination = config["destination"] @@ -275,9 +372,12 @@ def status_view() -> dict[str, Any]: "recovery_key_set": config["recovery_key"] is not None, "recovery_key_confirmed": bool(config["confirmed_recovery_key"]), "retention": config["retention"], + "offload": config["offload"], "schedule": config["schedule"], "last_backup": config["last_backup"], "last_prune": config["last_prune"], + "last_offload": config["last_offload"], + "last_verification": config["last_verification"], "hosted": hosted, } @@ -285,16 +385,22 @@ def status_view() -> dict[str, Any]: __all__ = [ "BACKUP_DEFAULTS", "BackupKeys", + "OFFLOAD_KEYS", + "OFFLOAD_STATUSES", + "VERIFICATION_STATUSES", "clear_backup_config", "generate_and_store_keys", "get_backup_config", "get_destination", "get_keys", "record_backup_result", + "record_offload_result", "record_prune_result", + "record_verification_result", "set_destination", "set_enabled", "set_mode", + "set_offload", "set_recovery_key", "set_recovery_key_confirmed", "set_retention", diff --git a/solstone/think/offload_ledger.py b/solstone/think/offload_ledger.py new file mode 100644 index 000000000..7d1ae4ead --- /dev/null +++ b/solstone/think/offload_ledger.py @@ -0,0 +1,496 @@ +# SPDX-License-Identifier: AGPL-3.0-only +# Copyright (c) 2026 sol pbc + +"""Durable media-offload ledger. + +This ledger is JSONL, not sqlite, because backup excludes ``*.sqlite*`` while +``health/`` is included. Losing the journal must not lose the map to the +owner's only backed-up copy of their media. It is fsync'd even though the +nearby pruning audit writer is not: pruning audits are post-hoc records, while +this ledger is a pre-delete witness for a later path that deletes local media +immediately after append returns. + +Per-file records intentionally use ``name`` and ``bytes`` to match the existing +raw-media audit shape, but use ``sha256`` instead of ``hash`` because this +durable schema rides encrypted backups and may be read years later by a restore +flow. A future delete path should compute each digest once, feed ``sha256`` to +this ledger, and remap only that key to the audit writer's ``hash`` at the call +site; it must not hash twice. + +Restore events carry only the identity spine and time. Repeating old file or +snapshot facts on restore would invite ambiguous folded reads; restore simply +invalidates the current offload state for a segment. + +Fold order is append order. Timestamps are informational, so a clock step +backward must not change the winning event. Read degradation is also deliberate: +unlike the repo's usual fail-loudly rule, a crashing read could let a future +teardown gate lose the owner's only media copy. A degraded zero is not a clean +zero; teardown gates must check ``degraded is False`` before trusting zero +offloaded bytes or segments. +""" + +from __future__ import annotations + +import json +import logging +import re +import time as time_module +from dataclasses import dataclass +from pathlib import Path +from typing import Any, Literal + +from solstone.think.journal_io import append_jsonl +from solstone.think.utils import ( + DATE_RE, + DEFAULT_STREAM, + STREAM_RE, + get_journal, + segment_key, +) + +logger = logging.getLogger(__name__) + +EVENT_OFFLOAD = "offload" +EVENT_RESTORE = "restore" +SHA256_RE = re.compile(r"[0-9a-f]{64}") + + +@dataclass(frozen=True) +class OffloadFile: + name: str + bytes: int + sha256: str + + +@dataclass(frozen=True) +class SegmentOffloadSummary: + day: str + stream: str + segment: str + currently_offloaded: bool + snapshot_id: str | None + files: tuple[OffloadFile, ...] + offloaded_bytes: int + offloaded_file_count: int + skipped_records: int + unreadable_ledgers: tuple[str, ...] + + @property + def degraded(self) -> bool: + return bool(self.unreadable_ledgers) + + +@dataclass(frozen=True) +class DayOffloadSummary: + day: str + segments: tuple[SegmentOffloadSummary, ...] + offloaded_bytes: int + offloaded_file_count: int + offloaded_segments: int + skipped_records: int + unreadable_ledgers: tuple[str, ...] + + @property + def degraded(self) -> bool: + return bool(self.unreadable_ledgers) + + +@dataclass(frozen=True) +class JournalOffloadSummary: + days: tuple[DayOffloadSummary, ...] + offloaded_bytes: int + offloaded_file_count: int + offloaded_segments: int + offloaded_days: int + skipped_records: int + unreadable_ledgers: tuple[str, ...] + + @property + def degraded(self) -> bool: + return bool(self.unreadable_ledgers) + + +@dataclass(frozen=True) +class _LedgerEvent: + event_kind: Literal["offload", "restore"] + time: int + day: str + stream: str + segment: str + snapshot_id: str | None + files: tuple[OffloadFile, ...] + + +@dataclass(frozen=True) +class _LedgerRead: + events: tuple[_LedgerEvent, ...] + skipped_records: int + unreadable_ledgers: tuple[str, ...] + + +def append_offload_event( + *, + day: str, + stream: str, + segment: str, + snapshot_id: str, + files: tuple[OffloadFile, ...] | list[OffloadFile], + time: int | None = None, +) -> None: + event_time = _write_event_time(time) + _validate_identity(day, stream, segment) + if not isinstance(snapshot_id, str) or not snapshot_id: + raise ValueError("snapshot_id must be a non-empty string") + if not files: + raise ValueError("files must not be empty") + validated_files = tuple(_validate_file_record(file) for file in files) + + append_jsonl( + _ledger_path(day), + { + "event_kind": EVENT_OFFLOAD, + "time": event_time, + "day": day, + "stream": stream, + "segment": segment, + "snapshot_id": snapshot_id, + "files": [_file_to_record(file) for file in validated_files], + }, + ) + + +def append_restore_event( + *, + day: str, + stream: str, + segment: str, + time: int | None = None, +) -> None: + event_time = _write_event_time(time) + _validate_identity(day, stream, segment) + + append_jsonl( + _ledger_path(day), + { + "event_kind": EVENT_RESTORE, + "time": event_time, + "day": day, + "stream": stream, + "segment": segment, + }, + ) + + +def summarize_segment(day: str, stream: str, segment: str) -> SegmentOffloadSummary: + _validate_identity(day, stream, segment) + read = _read_ledger_day(day) + states = _fold_events(read.events) + current = states.get((day, stream, segment)) + return _segment_summary( + day, + stream, + segment, + current, + skipped_records=read.skipped_records, + unreadable_ledgers=read.unreadable_ledgers, + ) + + +def summarize_day(day: str) -> DayOffloadSummary: + if not DATE_RE.fullmatch(day): + raise ValueError("day must be in YYYYMMDD format") + read = _read_ledger_day(day) + states = _fold_events(read.events) + segments = tuple( + _segment_summary( + key_day, + stream, + segment, + current, + skipped_records=0, + unreadable_ledgers=(), + ) + for (key_day, stream, segment), current in sorted(states.items()) + if key_day == day + ) + return _day_summary( + day, + segments, + skipped_records=read.skipped_records, + unreadable_ledgers=read.unreadable_ledgers, + ) + + +def summarize_journal() -> JournalOffloadSummary: + ledger_dir = _ledger_dir() + if not ledger_dir.is_dir(): + return JournalOffloadSummary( + days=(), + offloaded_bytes=0, + offloaded_file_count=0, + offloaded_segments=0, + offloaded_days=0, + skipped_records=0, + unreadable_ledgers=(), + ) + + days = tuple( + summarize_day(path.stem) + for path in sorted(ledger_dir.glob("*.jsonl")) + if DATE_RE.fullmatch(path.stem) + ) + unreadable_ledgers = tuple( + ledger for day in days for ledger in day.unreadable_ledgers + ) + return JournalOffloadSummary( + days=days, + offloaded_bytes=sum(day.offloaded_bytes for day in days), + offloaded_file_count=sum(day.offloaded_file_count for day in days), + offloaded_segments=sum(day.offloaded_segments for day in days), + offloaded_days=sum(1 for day in days if day.offloaded_segments > 0), + skipped_records=sum(day.skipped_records for day in days), + unreadable_ledgers=unreadable_ledgers, + ) + + +def ledger_path_for_day(day: str) -> Path: + if not DATE_RE.fullmatch(day): + raise ValueError("day must be in YYYYMMDD format") + return _ledger_path(day) + + +def _ledger_dir() -> Path: + return Path(get_journal()) / "health" / "offload" + + +def _ledger_path(day: str) -> Path: + return _ledger_dir() / f"{day}.jsonl" + + +def _write_event_time(value: int | None) -> int: + if value is None: + return int(time_module.time()) + return _read_event_time(value) + + +def _read_event_time(value: Any) -> int: + if type(value) is not int or value < 0: + raise ValueError("time must be a non-negative integer epoch second") + return value + + +def _validate_identity(day: str, stream: str, segment: str) -> None: + if not isinstance(day, str) or not DATE_RE.fullmatch(day): + raise ValueError("day must be in YYYYMMDD format") + if not isinstance(stream, str) or not ( + stream == DEFAULT_STREAM or STREAM_RE.fullmatch(stream) + ): + raise ValueError("stream must be a valid stream name") + if not isinstance(segment, str) or segment_key(segment) != segment: + raise ValueError("segment must be an exact segment key") + + +def _validate_file_record(file: OffloadFile) -> OffloadFile: + if not isinstance(file, OffloadFile): + raise ValueError("files must contain OffloadFile records") + _validate_file_fields(file.name, file.bytes, file.sha256) + return file + + +def _validate_file_fields(name: Any, size: Any, sha256: Any) -> None: + if not isinstance(name, str) or not name or "/" in name or name in {".", ".."}: + raise ValueError("file name must be a segment-local basename") + if type(size) is not int or size < 0: + raise ValueError("file bytes must be a non-negative integer") + if not isinstance(sha256, str) or SHA256_RE.fullmatch(sha256) is None: + raise ValueError("file sha256 must be a lowercase 64-character hex digest") + + +def _file_to_record(file: OffloadFile) -> dict[str, Any]: + return {"name": file.name, "bytes": file.bytes, "sha256": file.sha256} + + +def _read_ledger_day(day: str) -> _LedgerRead: + return _read_ledger_file(_ledger_path(day), expected_day=day) + + +def _read_ledger_file(path: Path, *, expected_day: str | None) -> _LedgerRead: + if not path.exists(): + return _LedgerRead(events=(), skipped_records=0, unreadable_ledgers=()) + + try: + raw = path.read_text(encoding="utf-8") + except (OSError, UnicodeDecodeError) as exc: + logger.warning("offload ledger read degraded for %s: %s", path, exc) + return _LedgerRead( + events=(), skipped_records=0, unreadable_ledgers=(str(path),) + ) + + events: list[_LedgerEvent] = [] + skipped_records = 0 + for lineno, line in enumerate(raw.splitlines(), start=1): + if not line.strip(): + continue + try: + record = json.loads(line) + event = _parse_event(record) + if expected_day is not None and event.day != expected_day: + raise ValueError("record day does not match ledger file day") + except json.JSONDecodeError: + skipped_records += 1 + logger.warning( + "skipping malformed offload ledger line %d in %s", lineno, path + ) + continue + except (TypeError, ValueError) as exc: + skipped_records += 1 + logger.warning( + "skipping invalid offload ledger record line %d in %s: %s", + lineno, + path, + exc, + ) + continue + events.append(event) + return _LedgerRead( + events=tuple(events), skipped_records=skipped_records, unreadable_ledgers=() + ) + + +def _parse_event(record: Any) -> _LedgerEvent: + if not isinstance(record, dict): + raise ValueError("record must be a JSON object") + event_kind = record.get("event_kind") + if event_kind == EVENT_OFFLOAD: + expected = { + "event_kind", + "time", + "day", + "stream", + "segment", + "snapshot_id", + "files", + } + if set(record) != expected: + raise ValueError("offload record has unexpected fields") + _validate_identity(record["day"], record["stream"], record["segment"]) + event_time = _read_event_time(record["time"]) + snapshot_id = record["snapshot_id"] + if not isinstance(snapshot_id, str) or not snapshot_id: + raise ValueError("snapshot_id must be a non-empty string") + raw_files = record["files"] + if not isinstance(raw_files, list) or not raw_files: + raise ValueError("files must be a non-empty list") + files = tuple(_parse_file_record(file) for file in raw_files) + return _LedgerEvent( + event_kind=EVENT_OFFLOAD, + time=event_time, + day=record["day"], + stream=record["stream"], + segment=record["segment"], + snapshot_id=snapshot_id, + files=files, + ) + + if event_kind == EVENT_RESTORE: + expected = {"event_kind", "time", "day", "stream", "segment"} + if set(record) != expected: + raise ValueError("restore record has unexpected fields") + _validate_identity(record["day"], record["stream"], record["segment"]) + return _LedgerEvent( + event_kind=EVENT_RESTORE, + time=_read_event_time(record["time"]), + day=record["day"], + stream=record["stream"], + segment=record["segment"], + snapshot_id=None, + files=(), + ) + + raise ValueError("event_kind must be offload or restore") + + +def _parse_file_record(record: Any) -> OffloadFile: + if not isinstance(record, dict) or set(record) != {"name", "bytes", "sha256"}: + raise ValueError("file record must contain name, bytes, sha256") + _validate_file_fields(record["name"], record["bytes"], record["sha256"]) + return OffloadFile( + name=record["name"], + bytes=record["bytes"], + sha256=record["sha256"], + ) + + +def _fold_events( + events: tuple[_LedgerEvent, ...], +) -> dict[tuple[str, str, str], _LedgerEvent]: + states: dict[tuple[str, str, str], _LedgerEvent] = {} + for event in events: + key = (event.day, event.stream, event.segment) + if event.event_kind == EVENT_OFFLOAD: + states[key] = event + elif event.event_kind == EVENT_RESTORE: + states.pop(key, None) + else: # pragma: no cover - parser owns the closed vocabulary. + raise RuntimeError(f"unexpected offload event kind: {event.event_kind}") + return states + + +def _segment_summary( + day: str, + stream: str, + segment: str, + current: _LedgerEvent | None, + *, + skipped_records: int, + unreadable_ledgers: tuple[str, ...], +) -> SegmentOffloadSummary: + files = current.files if current is not None else () + return SegmentOffloadSummary( + day=day, + stream=stream, + segment=segment, + currently_offloaded=current is not None, + snapshot_id=current.snapshot_id if current is not None else None, + files=files, + offloaded_bytes=sum(file.bytes for file in files), + offloaded_file_count=len(files), + skipped_records=skipped_records, + unreadable_ledgers=unreadable_ledgers, + ) + + +def _day_summary( + day: str, + segments: tuple[SegmentOffloadSummary, ...], + *, + skipped_records: int, + unreadable_ledgers: tuple[str, ...], +) -> DayOffloadSummary: + return DayOffloadSummary( + day=day, + segments=segments, + offloaded_bytes=sum(segment.offloaded_bytes for segment in segments), + offloaded_file_count=sum(segment.offloaded_file_count for segment in segments), + offloaded_segments=sum( + 1 for segment in segments if segment.currently_offloaded + ), + skipped_records=skipped_records, + unreadable_ledgers=unreadable_ledgers, + ) + + +__all__ = [ + "DayOffloadSummary", + "EVENT_OFFLOAD", + "EVENT_RESTORE", + "JournalOffloadSummary", + "OffloadFile", + "SegmentOffloadSummary", + "append_offload_event", + "append_restore_event", + "ledger_path_for_day", + "summarize_day", + "summarize_journal", + "summarize_segment", +] diff --git a/solstone/think/offload_measurement.py b/solstone/think/offload_measurement.py new file mode 100644 index 000000000..08efb5ccf --- /dev/null +++ b/solstone/think/offload_measurement.py @@ -0,0 +1,85 @@ +# SPDX-License-Identifier: AGPL-3.0-only +# Copyright (c) 2026 sol pbc + +"""Read-only measurement helpers for media offload.""" + +from __future__ import annotations + +import shutil +from dataclasses import dataclass +from pathlib import Path + +from solstone.think.retention import get_raw_media_files +from solstone.think.utils import day_dirs, get_journal, iter_segments + +MIN_FLOOR_BYTES = 20_000_000_000 + + +@dataclass(frozen=True) +class RawMediaDayUsage: + day: str + bytes: int + files: int + + +@dataclass(frozen=True) +class RawMediaUsage: + total_bytes: int + total_files: int + per_day: tuple[RawMediaDayUsage, ...] + + +@dataclass(frozen=True) +class SuggestedOffloadDefaults: + budget_bytes: int + floor_bytes: int + + +def measure_raw_media_usage() -> RawMediaUsage: + per_day: list[RawMediaDayUsage] = [] + total_bytes = 0 + total_files = 0 + + for day in sorted(day_dirs().keys()): + day_bytes = 0 + day_files = 0 + for _stream, _segment, segment_path in iter_segments(day): + for raw_file in get_raw_media_files(segment_path): + try: + size = raw_file.stat().st_size + except FileNotFoundError: + continue + day_bytes += size + day_files += 1 + per_day.append(RawMediaDayUsage(day=day, bytes=day_bytes, files=day_files)) + total_bytes += day_bytes + total_files += day_files + + return RawMediaUsage( + total_bytes=total_bytes, + total_files=total_files, + per_day=tuple(per_day), + ) + + +def device_free_bytes() -> int: + return shutil.disk_usage(Path(get_journal())).free + + +def suggest_offload_defaults(total_bytes: int) -> SuggestedOffloadDefaults: + if type(total_bytes) is not int or total_bytes <= 0: + raise ValueError("total_bytes must be a positive integer") + budget = total_bytes // 2 + floor = min(max(total_bytes // 10, MIN_FLOOR_BYTES), total_bytes // 4) + return SuggestedOffloadDefaults(budget_bytes=budget, floor_bytes=floor) + + +__all__ = [ + "MIN_FLOOR_BYTES", + "RawMediaDayUsage", + "RawMediaUsage", + "SuggestedOffloadDefaults", + "device_free_bytes", + "measure_raw_media_usage", + "suggest_offload_defaults", +] diff --git a/tests/test_backup_offload_guards.py b/tests/test_backup_offload_guards.py new file mode 100644 index 000000000..e016ca011 --- /dev/null +++ b/tests/test_backup_offload_guards.py @@ -0,0 +1,74 @@ +# SPDX-License-Identifier: AGPL-3.0-only +# Copyright (c) 2026 sol pbc + +from __future__ import annotations + +from fnmatch import fnmatchcase + +import pytest + +from solstone.think.backup.engine import BACKUP_EXCLUDES + + +def _restic_excludes(pattern: str, rel_path: str) -> bool: + path_parts = tuple(part for part in rel_path.strip("/").split("/") if part) + if "/" not in pattern: + return any(fnmatchcase(part, pattern) for part in path_parts) + + anchored = pattern.startswith("/") + pattern_parts = tuple(part for part in pattern.strip("/").split("/") if part) + starts = (0,) if anchored else range(len(path_parts)) + for start in starts: + for end in range(start, len(path_parts) + 1): + if _match_components(pattern_parts, path_parts[start:end]): + return True + return False + + +def _match_components( + pattern_parts: tuple[str, ...], path_parts: tuple[str, ...] +) -> bool: + if not pattern_parts: + return not path_parts + head, *tail = pattern_parts + rest = tuple(tail) + if head == "**": + return _match_components(rest, path_parts) or ( + bool(path_parts) and _match_components(pattern_parts, path_parts[1:]) + ) + return ( + bool(path_parts) + and fnmatchcase(path_parts[0], head) + and _match_components(rest, path_parts[1:]) + ) + + +@pytest.mark.parametrize( + ("pattern", "rel_path", "expected"), + [ + ("health", "health/offload/20260101.jsonl", True), + ("*.jsonl", "health/offload/20260101.jsonl", True), + ("health/*.jsonl", "health/offload/20260101.jsonl", False), + ("offload/*.jsonl", "health/offload/20260101.jsonl", True), + ("/health/offload", "health/offload/20260101.jsonl", True), + ("/offload/*.jsonl", "health/offload/20260101.jsonl", False), + ("health/**/20260101.jsonl", "health/offload/20260101.jsonl", True), + ], +) +def test_restic_exclude_component_model( + pattern: str, rel_path: str, expected: bool +) -> None: + assert _restic_excludes(pattern, rel_path) is expected + + +def test_offload_ledger_survives_backup_excludes() -> None: + # engine.py:58-62 documents why this exists: a bare "health" pattern + # previously matched health/ content by basename at every depth. + ledger_rel_path = "health/offload/20260101.jsonl" + + assert [ + pattern + for pattern in BACKUP_EXCLUDES + if _restic_excludes(pattern, ledger_rel_path) + ] == [] + assert _restic_excludes("health", ledger_rel_path) diff --git a/tests/test_backup_offload_state.py b/tests/test_backup_offload_state.py new file mode 100644 index 000000000..85e4e8dbb --- /dev/null +++ b/tests/test_backup_offload_state.py @@ -0,0 +1,187 @@ +# SPDX-License-Identifier: AGPL-3.0-only +# Copyright (c) 2026 sol pbc + +from __future__ import annotations + +import json +from pathlib import Path + +import pytest + +from solstone.think.backup import state + + +def _config_path(journal: Path) -> Path: + return journal / "config" / "journal.json" + + +def _write_config(journal: Path, payload: dict) -> None: + config_path = _config_path(journal) + config_path.parent.mkdir(parents=True, exist_ok=True) + config_path.write_text(json.dumps(payload, indent=2) + "\n", encoding="utf-8") + + +def _read_config(journal: Path) -> dict: + return json.loads(_config_path(journal).read_text(encoding="utf-8")) + + +def test_missing_offload_key_reads_pinned_defaults( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch +) -> None: + monkeypatch.setenv("SOLSTONE_JOURNAL", str(tmp_path)) + _write_config(tmp_path, {"backup": {"enabled": True}}) + + assert state.get_backup_config()["offload"] == { + "enabled": False, + "budget_bytes": None, + "floor_bytes": None, + } + + +def test_set_offload_valid_write_round_trips( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch +) -> None: + monkeypatch.setenv("SOLSTONE_JOURNAL", str(tmp_path)) + _write_config(tmp_path, {}) + offload = { + "enabled": True, + "budget_bytes": 500_000_000_000, + "floor_bytes": None, + } + + state.set_offload(offload) + + assert _read_config(tmp_path)["backup"]["offload"] == offload + assert state.get_backup_config()["offload"] == offload + + +@pytest.mark.parametrize( + "offload", + [ + {"enabled": 1, "budget_bytes": None, "floor_bytes": None}, + {"enabled": "true", "budget_bytes": None, "floor_bytes": None}, + {"enabled": True, "budget_bytes": True, "floor_bytes": None}, + {"enabled": True, "budget_bytes": False, "floor_bytes": None}, + {"enabled": True, "budget_bytes": 0, "floor_bytes": None}, + {"enabled": True, "budget_bytes": -1, "floor_bytes": None}, + {"enabled": True, "budget_bytes": "1", "floor_bytes": None}, + {"enabled": True, "budget_bytes": None, "floor_bytes": 0}, + {"enabled": True, "budget_bytes": None, "floor_bytes": -1}, + {"enabled": True, "budget_bytes": None, "floor_bytes": "1"}, + {"enabled": True, "budget_bytes": None}, + { + "enabled": True, + "budget_bytes": None, + "floor_bytes": None, + "extra": 1, + }, + ], +) +def test_set_offload_rejects_invalid_shapes_without_writing( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch, offload: dict +) -> None: + monkeypatch.setenv("SOLSTONE_JOURNAL", str(tmp_path)) + original = {"backup": {"enabled": False}} + _write_config(tmp_path, original) + before = _read_config(tmp_path) + + with pytest.raises(ValueError): + state.set_offload(offload) + + assert _read_config(tmp_path) == before + + +def test_set_offload_preserves_existing_backup_section_without_materializing_defaults( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch +) -> None: + monkeypatch.setenv("SOLSTONE_JOURNAL", str(tmp_path)) + backup = { + "enabled": False, + "mode": "operated", + "retention": {"hourly": 2, "daily": 3, "weekly": 4, "monthly": 5}, + "custom_marker": {"nested": ["keep", 7]}, + } + _write_config(tmp_path, {"backup": backup}) + before_backup = _read_config(tmp_path)["backup"] + + state.set_offload( + { + "enabled": True, + "budget_bytes": 10, + "floor_bytes": 5, + } + ) + + after_backup = _read_config(tmp_path)["backup"] + assert set(after_backup) == {*before_backup, "offload"} + for key, value in before_backup.items(): + assert json.dumps(after_backup[key], sort_keys=True) == json.dumps( + value, sort_keys=True + ) + assert "destination" not in after_backup + assert "last_prune" not in after_backup + + +def test_set_offload_has_no_cross_field_backup_enabled_rule( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch +) -> None: + monkeypatch.setenv("SOLSTONE_JOURNAL", str(tmp_path)) + _write_config(tmp_path, {"backup": {"enabled": False}}) + + state.set_offload( + { + "enabled": True, + "budget_bytes": 10, + "floor_bytes": 20, + } + ) + + assert _read_config(tmp_path)["backup"]["enabled"] is False + assert _read_config(tmp_path)["backup"]["offload"] == { + "enabled": True, + "budget_bytes": 10, + "floor_bytes": 20, + } + + +def test_offload_state_records_and_status_view( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch +) -> None: + monkeypatch.setenv("SOLSTONE_JOURNAL", str(tmp_path)) + _write_config(tmp_path, {"backup": {}}) + + config = state.get_backup_config() + assert config["last_offload"] == {"time": None, "status": None, "reason": None} + assert config["last_verification"] == { + "time": None, + "status": None, + "reason": None, + } + + state.record_offload_result(status="stalled", time=123, reason="no_progress") + state.record_verification_result(status="skipped", time=None, reason="disabled") + + backup = _read_config(tmp_path)["backup"] + assert backup["last_offload"] == { + "time": 123, + "status": "stalled", + "reason": "no_progress", + } + assert backup["last_verification"] == { + "time": None, + "status": "skipped", + "reason": "disabled", + } + with pytest.raises(ValueError): + state.record_offload_result(status="degraded", time=1) + with pytest.raises(ValueError): + state.record_verification_result(status="stalled", time=1) + + view = state.status_view() + assert view["offload"] == { + "enabled": False, + "budget_bytes": None, + "floor_bytes": None, + } + assert view["last_offload"] == backup["last_offload"] + assert view["last_verification"] == backup["last_verification"] diff --git a/tests/test_offload_ledger.py b/tests/test_offload_ledger.py new file mode 100644 index 000000000..92a2107c3 --- /dev/null +++ b/tests/test_offload_ledger.py @@ -0,0 +1,390 @@ +# SPDX-License-Identifier: AGPL-3.0-only +# Copyright (c) 2026 sol pbc + +from __future__ import annotations + +import json +import logging +import os +import stat +from collections.abc import Callable +from pathlib import Path + +import pytest + +from solstone.think.offload_ledger import ( + EVENT_OFFLOAD, + EVENT_RESTORE, + OffloadFile, + append_offload_event, + append_restore_event, + ledger_path_for_day, + summarize_day, + summarize_journal, + summarize_segment, +) + +DAY = "20260101" +STREAM = "archon" +SEGMENT = "120000_300" +SHA_A = "a" * 64 +SHA_B = "b" * 64 + + +def _is_ledger_fd_fsynced( + monkeypatch: pytest.MonkeyPatch, ledger_path: Path, writer: Callable[[], None] +) -> bool: + calls = [] + real_fsync = os.fsync + + def spy(fd: int) -> None: + fd_stat = os.fstat(fd) + try: + ledger_stat = ledger_path.stat() + except FileNotFoundError: + ledger_stat = None + if ( + ledger_stat is not None + and stat.S_ISREG(fd_stat.st_mode) + and (fd_stat.st_dev, fd_stat.st_ino) + == (ledger_stat.st_dev, ledger_stat.st_ino) + ): + calls.append(fd) + real_fsync(fd) + + monkeypatch.setattr("solstone.think.journal_io.append.os.fsync", spy) + writer() + return bool(calls) + + +def _write_plain_jsonl(path: Path, record: dict) -> None: + path.parent.mkdir(parents=True, exist_ok=True) + with open(path, "ab") as handle: + handle.write((json.dumps(record) + "\n").encode("utf-8")) + handle.flush() + + +def _offload_record( + *, + day: str, + stream: str, + segment: str, + snapshot_id: str, + size: int, + name: str = "audio.wav", + sha256: str = SHA_A, + time: int = 1, +) -> dict: + return { + "event_kind": EVENT_OFFLOAD, + "time": time, + "day": day, + "stream": stream, + "segment": segment, + "snapshot_id": snapshot_id, + "files": [{"name": name, "bytes": size, "sha256": sha256}], + } + + +def _restore_record( + *, + day: str, + stream: str, + segment: str, + time: int = 1, +) -> dict: + return { + "event_kind": EVENT_RESTORE, + "time": time, + "day": day, + "stream": stream, + "segment": segment, + } + + +def test_offload_append_writes_media_day_ledger_and_segment_summary( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch +) -> None: + monkeypatch.setenv("SOLSTONE_JOURNAL", str(tmp_path)) + + append_offload_event( + day=DAY, + stream=STREAM, + segment=SEGMENT, + snapshot_id="snap-1", + files=[ + OffloadFile(name="audio.wav", bytes=10, sha256=SHA_A), + OffloadFile(name="screen.mp4", bytes=20, sha256=SHA_B), + ], + time=10, + ) + + ledger_path = ledger_path_for_day(DAY) + assert ledger_path == tmp_path / "health" / "offload" / f"{DAY}.jsonl" + record = json.loads(ledger_path.read_text(encoding="utf-8").splitlines()[0]) + assert "schema_version" not in record + assert record["event_kind"] == EVENT_OFFLOAD + assert record["snapshot_id"] == "snap-1" + assert record["files"] == [ + {"name": "audio.wav", "bytes": 10, "sha256": SHA_A}, + {"name": "screen.mp4", "bytes": 20, "sha256": SHA_B}, + ] + assert "hash" not in record["files"][0] + + summary = summarize_segment(DAY, STREAM, SEGMENT) + assert summary.currently_offloaded is True + assert summary.snapshot_id == "snap-1" + assert summary.offloaded_bytes == 30 + assert summary.offloaded_file_count == 2 + assert summary.files[0].sha256 == SHA_A + + +def test_restore_and_reoffload_fold_by_append_order_not_time( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch +) -> None: + monkeypatch.setenv("SOLSTONE_JOURNAL", str(tmp_path)) + + append_offload_event( + day=DAY, + stream=STREAM, + segment=SEGMENT, + snapshot_id="snap-old", + files=[OffloadFile(name="audio.wav", bytes=10, sha256=SHA_A)], + time=300, + ) + append_restore_event(day=DAY, stream=STREAM, segment=SEGMENT, time=200) + + restored = summarize_segment(DAY, STREAM, SEGMENT) + assert restored.currently_offloaded is False + assert restored.offloaded_bytes == 0 + assert restored.snapshot_id is None + + append_offload_event( + day=DAY, + stream=STREAM, + segment=SEGMENT, + snapshot_id="snap-new", + files=[OffloadFile(name="audio.wav", bytes=11, sha256=SHA_B)], + time=100, + ) + + reoffloaded = summarize_segment(DAY, STREAM, SEGMENT) + assert reoffloaded.currently_offloaded is True + assert reoffloaded.snapshot_id == "snap-new" + assert reoffloaded.offloaded_bytes == 11 + + +def test_append_fsyncs_ledger_fd_before_return( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch +) -> None: + monkeypatch.setenv("SOLSTONE_JOURNAL", str(tmp_path)) + ledger_path = ledger_path_for_day(DAY) + + with monkeypatch.context() as context: + assert _is_ledger_fd_fsynced( + context, + ledger_path, + lambda: append_offload_event( + day=DAY, + stream=STREAM, + segment=SEGMENT, + snapshot_id="snap-1", + files=[OffloadFile(name="audio.wav", bytes=10, sha256=SHA_A)], + time=1, + ), + ) + + plain_path = tmp_path / "health" / "offload" / "plain.jsonl" + with monkeypatch.context() as context: + assert not _is_ledger_fd_fsynced( + context, + plain_path, + lambda: _write_plain_jsonl( + plain_path, + _offload_record( + day=DAY, + stream=STREAM, + segment=SEGMENT, + snapshot_id="plain", + size=10, + ), + ), + ) + + +def test_day_and_journal_summaries_fold_mixed_days( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch +) -> None: + monkeypatch.setenv("SOLSTONE_JOURNAL", str(tmp_path)) + append_offload_event( + day="20260101", + stream=STREAM, + segment="120000_300", + snapshot_id="full", + files=[ + OffloadFile(name="a.wav", bytes=10, sha256=SHA_A), + OffloadFile(name="b.mp4", bytes=5, sha256=SHA_B), + ], + time=1, + ) + append_offload_event( + day="20260102", + stream=STREAM, + segment="120000_300", + snapshot_id="restored", + files=[OffloadFile(name="a.wav", bytes=20, sha256=SHA_A)], + time=1, + ) + append_restore_event(day="20260102", stream=STREAM, segment="120000_300", time=2) + append_offload_event( + day="20260103", + stream=STREAM, + segment="120000_300", + snapshot_id="mixed-kept", + files=[OffloadFile(name="a.wav", bytes=30, sha256=SHA_A)], + time=1, + ) + append_offload_event( + day="20260103", + stream=STREAM, + segment="121000_300", + snapshot_id="mixed-restored", + files=[OffloadFile(name="b.wav", bytes=40, sha256=SHA_B)], + time=1, + ) + append_restore_event(day="20260103", stream=STREAM, segment="121000_300", time=2) + + assert summarize_day("20260101").offloaded_bytes == 15 + assert summarize_day("20260102").offloaded_bytes == 0 + mixed = summarize_day("20260103") + assert mixed.offloaded_bytes == 30 + assert mixed.offloaded_segments == 1 + + journal = summarize_journal() + assert [day.day for day in journal.days] == ["20260101", "20260102", "20260103"] + assert journal.offloaded_bytes == 45 + assert journal.offloaded_file_count == 3 + assert journal.offloaded_segments == 2 + assert journal.offloaded_days == 2 + + +def test_malformed_line_warns_counts_skipped_and_keeps_valid_totals( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch, caplog: pytest.LogCaptureFixture +) -> None: + monkeypatch.setenv("SOLSTONE_JOURNAL", str(tmp_path)) + ledger_path = ledger_path_for_day(DAY) + ledger_path.parent.mkdir(parents=True) + ledger_path.write_text( + "\n".join( + [ + json.dumps( + _offload_record( + day=DAY, + stream=STREAM, + segment="120000_300", + snapshot_id="snap-1", + size=10, + ) + ), + "{bad", + json.dumps( + _offload_record( + day=DAY, + stream=STREAM, + segment="121000_300", + snapshot_id="snap-2", + size=20, + sha256=SHA_B, + ) + ), + ] + ) + + "\n", + encoding="utf-8", + ) + + with caplog.at_level(logging.WARNING, logger="solstone.think.offload_ledger"): + summary = summarize_day(DAY) + + assert summary.offloaded_bytes == 30 + assert summary.skipped_records == 1 + assert any(str(ledger_path) in record.message for record in caplog.records) + assert any("malformed" in record.message for record in caplog.records) + + +def test_null_time_record_is_skipped_not_fabricated( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch, caplog: pytest.LogCaptureFixture +) -> None: + monkeypatch.setenv("SOLSTONE_JOURNAL", str(tmp_path)) + ledger_path = ledger_path_for_day(DAY) + invalid = _offload_record( + day=DAY, + stream=STREAM, + segment="120000_300", + snapshot_id="bad-time", + size=20, + ) + invalid["time"] = None + ledger_path.parent.mkdir(parents=True) + ledger_path.write_text( + "\n".join( + [ + json.dumps( + _offload_record( + day=DAY, + stream=STREAM, + segment="121000_300", + snapshot_id="valid", + size=10, + ) + ), + json.dumps(invalid), + ] + ) + + "\n", + encoding="utf-8", + ) + + with caplog.at_level(logging.WARNING, logger="solstone.think.offload_ledger"): + summary = summarize_day(DAY) + + assert summary.offloaded_bytes == 10 + assert summary.skipped_records == 1 + assert any(str(ledger_path) in record.message for record in caplog.records) + + +def test_undecodable_ledger_degrades_without_clean_zero( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch, caplog: pytest.LogCaptureFixture +) -> None: + monkeypatch.setenv("SOLSTONE_JOURNAL", str(tmp_path)) + ledger_path = ledger_path_for_day(DAY) + ledger_path.parent.mkdir(parents=True) + ledger_path.write_bytes(b"\xff") + + with caplog.at_level(logging.WARNING, logger="solstone.think.offload_ledger"): + summary = summarize_day(DAY) + + assert summary.offloaded_bytes == 0 + assert summary.degraded is True + assert summary.unreadable_ledgers == (str(ledger_path),) + assert any(str(ledger_path) in record.message for record in caplog.records) + + +def test_absent_or_empty_ledger_is_clean_zero( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch +) -> None: + monkeypatch.setenv("SOLSTONE_JOURNAL", str(tmp_path)) + + missing = summarize_journal() + assert missing.offloaded_bytes == 0 + assert missing.skipped_records == 0 + assert missing.degraded is False + + ledger_path = ledger_path_for_day(DAY) + ledger_path.parent.mkdir(parents=True) + ledger_path.write_text("", encoding="utf-8") + + empty = summarize_day(DAY) + assert empty.offloaded_bytes == 0 + assert empty.skipped_records == 0 + assert empty.degraded is False diff --git a/tests/test_offload_measurement.py b/tests/test_offload_measurement.py new file mode 100644 index 000000000..c242330bf --- /dev/null +++ b/tests/test_offload_measurement.py @@ -0,0 +1,123 @@ +# SPDX-License-Identifier: AGPL-3.0-only +# Copyright (c) 2026 sol pbc + +from __future__ import annotations + +from pathlib import Path +from types import SimpleNamespace + +import pytest + +from solstone.think import offload_measurement +from solstone.think.offload_measurement import ( + device_free_bytes, + measure_raw_media_usage, + suggest_offload_defaults, +) +from solstone.think.retention import compute_storage_summary + +GB = 10**9 + + +def _segment(journal: Path, day: str, stream: str = "archon") -> Path: + path = journal / "chronicle" / day / stream / "120000_300" + path.mkdir(parents=True, exist_ok=True) + return path + + +def _write_bytes(path: Path, size: int) -> Path: + path.parent.mkdir(parents=True, exist_ok=True) + path.write_bytes(b"x" * size) + return path + + +def test_raw_media_measurement_matches_retention_predicate_non_recursive( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch +) -> None: + monkeypatch.setenv("SOLSTONE_JOURNAL", str(tmp_path)) + segment = _segment(tmp_path, "20260101") + _write_bytes(segment / "audio.wav", 10) + _write_bytes(segment / "clip.mp4", 20) + _write_bytes(segment / "frame.png", 30) + _write_bytes(segment / "monitor_0_diff.png", 40) + _write_bytes(segment / "audio.jsonl", 50) + _write_bytes(segment / "note.json", 60) + _write_bytes(segment / "talents" / "frame.png", 70) + + usage = measure_raw_media_usage() + + assert usage.total_bytes == 100 + assert usage.total_files == 4 + assert usage.per_day == (offload_measurement.RawMediaDayUsage("20260101", 100, 4),) + assert compute_storage_summary().raw_media_bytes == usage.total_bytes + + +def test_raw_media_per_day_breakdown_is_chronological( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch +) -> None: + monkeypatch.setenv("SOLSTONE_JOURNAL", str(tmp_path)) + _write_bytes(_segment(tmp_path, "20260301") / "audio.wav", 3) + _write_bytes(_segment(tmp_path, "20260101") / "audio.wav", 1) + _write_bytes(_segment(tmp_path, "20260201") / "audio.wav", 2) + + usage = measure_raw_media_usage() + + assert [day.day for day in usage.per_day] == [ + "20260101", + "20260201", + "20260301", + ] + assert [day.bytes for day in usage.per_day] == [1, 2, 3] + + +def test_raw_media_measurement_tolerates_file_vanishing_before_stat( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch +) -> None: + monkeypatch.setenv("SOLSTONE_JOURNAL", str(tmp_path)) + segment = _segment(tmp_path, "20260101") + vanishing = _write_bytes(segment / "vanishing.wav", 10) + stable = _write_bytes(segment / "stable.wav", 20) + + def fake_get_raw_media_files(_segment_path: Path) -> list[Path]: + vanishing.unlink() + return [vanishing, stable] + + monkeypatch.setattr( + offload_measurement, "get_raw_media_files", fake_get_raw_media_files + ) + + usage = measure_raw_media_usage() + + assert usage.total_bytes == 20 + assert usage.total_files == 1 + + +def test_device_free_bytes_uses_disk_usage_free( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch +) -> None: + monkeypatch.setenv("SOLSTONE_JOURNAL", str(tmp_path)) + + def fake_disk_usage(path: Path) -> SimpleNamespace: + assert path == tmp_path + return SimpleNamespace(total=1000 * GB, used=100 * GB, free=850 * GB) + + monkeypatch.setattr(offload_measurement.shutil, "disk_usage", fake_disk_usage) + + assert device_free_bytes() == 850 * GB + + +@pytest.mark.parametrize( + ("total", "budget", "floor"), + [ + (1000 * GB, 500 * GB, 100 * GB), + (100 * GB, 50 * GB, 20 * GB), + (30 * GB, 15 * GB, 7_500_000_000), + ], +) +def test_suggest_offload_defaults_decimal_gb( + total: int, budget: int, floor: int +) -> None: + defaults = suggest_offload_defaults(total) + + assert defaults.budget_bytes == budget + assert defaults.floor_bytes == floor