From 5a5f8b011c039a4f504a772f48d8756ef40b8b90 Mon Sep 17 00:00:00 2001 From: Leo Vasanko Date: Wed, 2 Sep 2026 16:24:26 +0000 Subject: [PATCH] rotation: implement retention-based history rotation New kanta.rotation module plans rotation from in-memory bytes: walks snapshots backwards to a replay base predating the cutoff, validates replayed state against every snapshot in range, and rebuilds the file as leading cutoff snapshot + re-framed retained changes (+ final snapshot when enough changes survived). KantaImpl.open executes the plan under the exclusive lock before replay: copy2 aside to {stem}@{ts}.kantadb, in-place locked rewrite with fdatasync, then trim the rotated copy to the dropped-history prefix. Enabled via Kanta(retention=timedelta). --- kanta/kanta.py | 10 ++- kanta/kantaimpl.py | 29 +++++- kanta/rotation.py | 220 +++++++++++++++++++++++++++++++++++++++++++++ kanta/snapshot.py | 5 ++ 4 files changed, 262 insertions(+), 2 deletions(-) create mode 100644 kanta/rotation.py diff --git a/kanta/kanta.py b/kanta/kanta.py index 7e01d91..30feac0 100644 --- a/kanta/kanta.py +++ b/kanta/kanta.py @@ -2,7 +2,7 @@ from __future__ import annotations import logging -from datetime import datetime +from datetime import datetime, timedelta from pathlib import Path from types import ModuleType, SimpleNamespace from typing import Any, Generic, TypeVar @@ -53,6 +53,7 @@ class Kanta(Generic[T]): migrations: ModuleType | str | None = None, serializer: Serializer | None = None, flush_interval: float = 0.1, + retention: timedelta | None = None, ): """Initialize a Kanta persistence instance. @@ -63,6 +64,12 @@ class Kanta(Generic[T]): migrations: Optional migrations module object or import path. flush_interval: Background flush interval in seconds. serializer: Optional serializer implementation. + retention: Optional history retention window. When set, opening the + database rotates it: history older than ``now - retention`` is + moved to a ``{stem}@{timestamp}.kantadb`` sibling file and the + main file is rewritten with a fresh snapshot plus the retained + records (see ``docs/rotation.md``). ``None`` (default) disables + rotation. Raises: ImportError: If ``migrations`` is a string path that cannot be imported. @@ -78,6 +85,7 @@ class Kanta(Generic[T]): type=data_type, migrations=migrations, flush_interval=flush_interval, + retention=retention, kanta=self, ) diff --git a/kanta/kantaimpl.py b/kanta/kantaimpl.py index 38b4d6e..220ef3d 100644 --- a/kanta/kantaimpl.py +++ b/kanta/kantaimpl.py @@ -6,7 +6,7 @@ import asyncio import copy import importlib import logging -from datetime import UTC, datetime +from datetime import UTC, datetime, timedelta from types import SimpleNamespace from typing import Any, Generic, TypeVar @@ -21,6 +21,7 @@ from kanta.logging import ( ) from kanta.migrations import MigrationReport, Migrations from kanta.persistence import PersistenceMixin +from kanta.rotation import execute_rotation, plan_rotation from kanta.serialization import restore_data_in_place, struct_to_dict from kanta.serialization.base import replay @@ -42,6 +43,7 @@ class KantaImpl(PersistenceMixin, Generic[T]): self.data: T = kwargs.pop("data") self._kanta = kwargs.pop("kanta", None) migrations = kwargs.pop("migrations", None) + self.retention: timedelta | None = kwargs.pop("retention", None) self.ctx = SimpleNamespace() super().__init__(**kwargs) self.migrations: Migrations | None = None @@ -181,6 +183,9 @@ class KantaImpl(PersistenceMixin, Generic[T]): # From this point the file is open and must be closed via close(). self.opened = True + if content and self.retention is not None and not readonly: + content = await self._maybe_rotate(content, log) + if content: try: rr = replay( @@ -375,6 +380,28 @@ class KantaImpl(PersistenceMixin, Generic[T]): if not self.readonly: self.background_task = asyncio.create_task(self._background_loop()) + async def _maybe_rotate(self, content: bytes, log: bool | logging.Logger) -> bytes: + """Rotate history older than the retention window (see docs/rotation.md). + + Runs while the file is locked and quiescent, before replay. Returns + the (possibly replaced) content to replay. Rotation failures abort the + open with the original file intact. + """ + cutoff = self.now() - self.retention + plan = await asyncio.to_thread( + plan_rotation, + content, + framer=self.framer, + serializer=self.serializer, + cutoff=cutoff, + now=self.now(), + min_diffs=self.snapshot.min_diffs, + ) + if plan is None: + return content + await asyncio.to_thread(execute_rotation, self.filename, self.file, plan, log=log) + return plan.new_content + async def close(self) -> None: """Stop the background task, flush pending changes, and release the file lock.""" if not self.opened: diff --git a/kanta/rotation.py b/kanta/rotation.py new file mode 100644 index 0000000..410e535 --- /dev/null +++ b/kanta/rotation.py @@ -0,0 +1,220 @@ +"""Database rotation: bound on-disk history to a retention window. + +See docs/rotation.md for the design. All planning happens on the in-memory +bytes of the database file; the caller (KantaImpl.open) performs the actual +copy-aside, in-place rewrite and rotated-file trimming under the file lock. +""" + +from __future__ import annotations + +import copy +import logging +import shutil +from dataclasses import dataclass +from datetime import datetime +from pathlib import Path +from typing import Any + +from kanta.exceptions import DatabaseError +from kanta.structs import ChangeRecord, Snapshot +from kanta.serialization.base import Serializer, apply_diff +from kanta.serialization.framing import Framer + +_logger = logging.getLogger(__name__) + + +@dataclass +class RotationPlan: + """Everything needed to execute a rotation on disk.""" + + new_content: bytes + cutoff_end: int # byte length of the dropped-history prefix + rotated_ts: datetime # ts of the last dropped change record + retained_changes: int + + +def rotated_path_for(path: Path, ts: datetime) -> Path: + """Sibling path for the rotated history: ``{stem}@{ISO-basic-ts}.kantadb``. + + The timestamp uses ISO 8601 basic format with microsecond precision so it + matches the record ``ts`` values stored in the database. On collision an + incrementing suffix is inserted before the extension. + """ + stamp = ts.strftime("%Y%m%dT%H%M%S.%fZ") + candidate = path.with_name(f"{path.stem}@{stamp}.kantadb") + n = 1 + while candidate.exists(): + candidate = path.with_name(f"{path.stem}@{stamp}.{n}.kantadb") + n += 1 + return candidate + + +class _Entry: + """One parsed record frame with its byte range.""" + + __slots__ = ("is_snapshot", "record", "byte_pos", "end_pos") + + def __init__( + self, + is_snapshot: bool, + record: ChangeRecord | Snapshot, + byte_pos: int, + end_pos: int, + ) -> None: + self.is_snapshot = is_snapshot + self.record = record + self.byte_pos = byte_pos + self.end_pos = end_pos + + +def _scan( + content: bytes, *, framer: Framer, serializer: Serializer +) -> list[_Entry]: + """Decode every record in *content* with byte ranges.""" + raw = list(framer.iter_records(content, 0)) + entries: list[_Entry] = [] + for i, (is_snapshot, payload, _line, byte_pos) in enumerate(raw): + end_pos = raw[i + 1][3] if i + 1 < len(raw) else len(content) + record = serializer.decode( + payload, type=Snapshot if is_snapshot else ChangeRecord + ) + entries.append(_Entry(is_snapshot, record, byte_pos, end_pos)) + return entries + + +def plan_rotation( + content: bytes, + *, + framer: Framer, + serializer: Serializer, + cutoff: datetime, + now: datetime, + min_diffs: int, +) -> RotationPlan | None: + """Plan a rotation of *content*, or return None when there is nothing to do. + + Raises: + DatabaseError: If replay from the chosen base snapshot does not match + a snapshot found inside the file (corrupt history). Rotation must + be aborted and the original file left untouched. + """ + if not content: + return None + entries = _scan(content, framer=framer, serializer=serializer) + changes = [e for e in entries if not e.is_snapshot] + if not changes: + return None # snapshot-only file: already fully rotated + + dropped = [e for e in changes if e.record.ts < cutoff] + if not dropped: + return None # retention window covers all history + + retained = [e for e in changes if e.record.ts >= cutoff] + rotated_ts = dropped[-1].record.ts + cutoff_end = retained[0].byte_pos if retained else len(content) + + # Replay base: walk snapshots newest-first and take the first (newest) + # one predating the cutoff; fall back to start of file. + base: _Entry | None = None + for e in reversed([e for e in entries if e.is_snapshot]): + if e.record.ts <= cutoff: + base = e + break + + state: dict[str, Any] = {} + version = 0 + m: datetime | None = None + if base is not None: + snap = base.record + assert isinstance(snap, Snapshot) + state = dict(snap.state) + version = snap.v + m = snap.m + + state_at_cutoff: dict[str, Any] | None = None + version_at_cutoff = version + m_at_cutoff = m + final_version = version + final_m = m + + for e in entries: + if base is not None and e.byte_pos <= base.byte_pos: + continue + if e.is_snapshot: + snap = e.record + assert isinstance(snap, Snapshot) + if snap.state != state: + raise DatabaseError( + "rotation aborted: replayed state does not match snapshot " + f"at byte {e.byte_pos}", + action="rotate", + ) + if snap.m is not None: + m = snap.m + continue + change = e.record + assert isinstance(change, ChangeRecord) + state = apply_diff(state, change.diff) + version = change.v + if change.m is not None: + m = change.m + if change.ts < cutoff: + state_at_cutoff = copy.deepcopy(state) + version_at_cutoff = version + m_at_cutoff = m + final_version = version + final_m = m + + # There is at least one dropped change, so the cutoff state is known. + assert state_at_cutoff is not None + + # Build the new content: leading cutoff snapshot, retained changes + # re-framed at fresh offsets, and a final snapshot when enough changes + # survived to warrant one (mirrors the regular snapshot policy). + out = bytearray() + leading = serializer.encode( + Snapshot(ts=rotated_ts, v=version_at_cutoff, state=state_at_cutoff, m=m_at_cutoff) + ) + out += framer.frame_snapshot(leading, record_offset=0) + for e in retained: + payload = serializer.encode(e.record) + out += framer.frame_change(payload, record_offset=len(out)) + if len(retained) >= min_diffs: + closing = serializer.encode( + Snapshot(ts=now, v=final_version, state=state, m=final_m) + ) + out += framer.frame_snapshot(closing, record_offset=len(out)) + + return RotationPlan( + new_content=bytes(out), + cutoff_end=cutoff_end, + rotated_ts=rotated_ts, + retained_changes=len(retained), + ) + + +def execute_rotation( + path: Path, file, plan: RotationPlan, *, log: bool | logging.Logger = True +) -> Path: + """Execute a planned rotation on disk. Caller must hold the lock on *file*. + + 1. Copy the original content aside to ``{stem}@{ts}.kantadb``. + 2. Rewrite the locked file in place with the new content and fsync. + 3. Trim the rotated copy to the dropped-history prefix. + + Returns the rotated file path. + """ + rotated = rotated_path_for(path, plan.rotated_ts) + shutil.copy2(path, rotated) + file.replace_content(plan.new_content) + with open(rotated, "r+b") as f: + f.truncate(plan.cutoff_end) + if log: + _logger.info( + "rotated %s: kept %d change record(s), history before %s moved to %s", + path, + plan.retained_changes, + plan.rotated_ts.isoformat(), + rotated, + ) + return rotated diff --git a/kanta/snapshot.py b/kanta/snapshot.py index e97357b..dd92ebe 100644 --- a/kanta/snapshot.py +++ b/kanta/snapshot.py @@ -35,6 +35,11 @@ class SnapshotState: """Force snapshot write on next check.""" self._force_pending = True + @property + def min_diffs(self) -> int: + """Minimum accumulated changes before a snapshot may be written.""" + return self._min_diffs + def record_changes(self, count: int) -> None: self.changes += count