Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
ab5b7584c9 | ||
|
|
d33f3f9c2f | ||
|
|
aee6c13996 |
@@ -51,100 +51,14 @@ asyncio.run(main())
|
||||
3. Let Kanta flush queued changes to disk in the background.
|
||||
4. Use snapshots and replay for fast startup and full history.
|
||||
|
||||
## Bootstrap and Open Modes
|
||||
## Documentation
|
||||
|
||||
When `open()` creates a brand-new database, it always writes a single bootstrap
|
||||
change record from the initial data object you passed to `Kanta(...)`. The
|
||||
simplest bootstrap is therefore the object itself — no extra code is required.
|
||||
- [Usage patterns](https://git.zi.fi/LeoVasanko/kanta/src/branch/main/docs/usage.md) — opening, data ownership, and lifecycle patterns
|
||||
- [Bootstrap and open modes](https://git.zi.fi/LeoVasanko/kanta/src/branch/main/docs/bootstrap.md) — seeding new databases, strict and read-only opens
|
||||
- [Validation](https://git.zi.fi/LeoVasanko/kanta/src/branch/main/docs/validation.md) — `@kanta.validate` integrity checks on open and transactions
|
||||
- [Migrations](https://git.zi.fi/LeoVasanko/kanta/src/branch/main/docs/migrations.md) — versioned schema evolution with `migrate_vN`
|
||||
- [Retention and rotation](https://git.zi.fi/LeoVasanko/kanta/src/branch/main/docs/rotation.md) — bounding history to a time window
|
||||
- [Fatal error handlers](https://git.zi.fi/LeoVasanko/kanta/src/branch/main/docs/fatal-errors.md) — observing background write failures
|
||||
- [On-disk format](https://git.zi.fi/LeoVasanko/kanta/src/branch/main/docs/database.md) — record layout and invariants
|
||||
|
||||
Bootstrap handlers are optional. Use them only when you need to modify the
|
||||
initial state at creation time, for example to seed defaults or perform
|
||||
expensive/external setup that should happen exactly once:
|
||||
|
||||
```python
|
||||
kanta = Kanta("data.kantadb", Data())
|
||||
|
||||
@kanta.bootstrap(action="seed", user="system")
|
||||
def seed_defaults(data) -> None:
|
||||
data.users["admin"] = User(name="Admin")
|
||||
|
||||
await kanta.open()
|
||||
```
|
||||
|
||||
You can also use `@kanta.bootstrap` with no arguments and async handlers:
|
||||
|
||||
```python
|
||||
@kanta.bootstrap
|
||||
async def bootstrap_async(data) -> None:
|
||||
data.counter = 1
|
||||
```
|
||||
|
||||
Whether or not handlers are registered, exactly one bootstrap change record is
|
||||
written when a new database is created. The record contains the initial object,
|
||||
or the state after all bootstrap handlers have run. When handlers are present:
|
||||
- they run in registration order,
|
||||
- bootstrap metadata (`action`, `user`, `mtime`) is taken from the last
|
||||
registration.
|
||||
|
||||
If any bootstrap handler raises, Kanta closes and removes the database file,
|
||||
then re-raises the error.
|
||||
|
||||
`open()` also supports strict open mode:
|
||||
|
||||
```python
|
||||
await kanta.open(create=False)
|
||||
```
|
||||
|
||||
With `create=False`, open fails if the database file does not exist or is
|
||||
empty.
|
||||
|
||||
Read-only mode opens an existing database without locking it or starting the
|
||||
background flush task. This is useful for readers that must not block the
|
||||
writer or modify the file:
|
||||
|
||||
```python
|
||||
await kanta.open(readonly=True)
|
||||
```
|
||||
|
||||
In read-only mode, records are replayed and migrations are applied in memory,
|
||||
but transactions and explicit flushes are rejected and the file is never
|
||||
created if missing.
|
||||
|
||||
## Fatal Error Handlers
|
||||
|
||||
Fatal background write errors can be observed with a decorator:
|
||||
|
||||
```python
|
||||
import os
|
||||
import signal
|
||||
|
||||
@kanta.fatal_error
|
||||
async def on_fatal(err):
|
||||
os.kill(os.getpid(), signal.SIGTERM) # Die
|
||||
```
|
||||
|
||||
Multiple fatal handlers are supported and run in registration order.
|
||||
|
||||
## Migrations
|
||||
|
||||
Adding or removing a field and other such simple operations are automatic, but when the time comes to really change your data model, implement a `migrate_v1` function that converts your old data to the new form. This works on plain built-in dict and other types, to avoid needing to preserve old versions of your structs.
|
||||
|
||||
Pass a module (or import path) containing `migrate_vN` functions:
|
||||
|
||||
```python
|
||||
kanta = Kanta("data.kantadb", Data(), migrations="myapp.migrations")
|
||||
await kanta.open()
|
||||
```
|
||||
|
||||
Kanta tracks migration version metadata automatically, and fast forwards your database to current version by running all the migrations needed while opening the database.
|
||||
|
||||
## On-Disk Format
|
||||
|
||||
Kanta in JSON mode (default) stores newline-delimited records. Transaction history is viewable by any simple text editor, and rollbacks to prior state are done by simply removing final lines (one per transaction)
|
||||
|
||||
MsgPack mode uses binary records with length and checksum to avoid data corruption.
|
||||
|
||||
- Change line: JSON object with metadata + `diff`
|
||||
- Snapshot line: `SNAPSHOT { ... full state ... }`
|
||||
|
||||
See `docs/database.md` for format details and invariants.
|
||||
Kanta in JSON mode (default) stores newline-delimited records, so transaction history is viewable in any text editor; MsgPack mode uses binary records with length and checksum to guard against corruption.
|
||||
|
||||
@@ -0,0 +1,48 @@
|
||||
# Bootstrap and Open Modes
|
||||
|
||||
When `open()` creates a brand-new database, it always writes a single bootstrap change record from the initial data object you passed to `Kanta(...)`. The simplest bootstrap is therefore the object itself — no extra code is required.
|
||||
|
||||
Bootstrap handlers are optional. Use them only when you need to modify the initial state at creation time, for example to seed defaults or perform expensive/external setup that should happen exactly once:
|
||||
|
||||
```python
|
||||
kanta = Kanta("data.kantadb", Data())
|
||||
|
||||
@kanta.bootstrap(action="seed", user="system")
|
||||
def seed_defaults(data) -> None:
|
||||
data.users["admin"] = User(name="Admin")
|
||||
|
||||
await kanta.open()
|
||||
```
|
||||
|
||||
You can also use `@kanta.bootstrap` with no arguments and async handlers:
|
||||
|
||||
```python
|
||||
@kanta.bootstrap
|
||||
async def bootstrap_async(data) -> None:
|
||||
data.counter = 1
|
||||
```
|
||||
|
||||
Whether or not handlers are registered, exactly one bootstrap change record is written when a new database is created. The record contains the initial object, or the state after all bootstrap handlers have run. When handlers are present:
|
||||
|
||||
- they run in registration order,
|
||||
- bootstrap metadata (`action`, `user`, `mtime`) is taken from the last registration.
|
||||
|
||||
If any bootstrap handler raises, Kanta closes and removes the database file, then re-raises the error.
|
||||
|
||||
## Strict open mode
|
||||
|
||||
```python
|
||||
await kanta.open(create=False)
|
||||
```
|
||||
|
||||
With `create=False`, open fails if the database file does not exist or is empty.
|
||||
|
||||
## Read-only mode
|
||||
|
||||
Read-only mode opens an existing database without locking it or starting the background flush task. This is useful for readers that must not block the writer or modify the file:
|
||||
|
||||
```python
|
||||
await kanta.open(readonly=True)
|
||||
```
|
||||
|
||||
In read-only mode, records are replayed and migrations are applied in memory, but transactions and explicit flushes are rejected and the file is never created if missing.
|
||||
+44
-140
@@ -1,13 +1,13 @@
|
||||
# Kanta Database Format and Design Principles
|
||||
|
||||
This document describes the on-disk format and design principles of Kanta.
|
||||
It is intentionally focused on the current standalone package behavior.
|
||||
This document describes the on-disk format and design principles of Kanta. It is intentionally focused on the current standalone package behavior.
|
||||
|
||||
## Core Principles
|
||||
|
||||
1. Append-only durability
|
||||
- State changes are persisted as appended JSON lines.
|
||||
- Existing lines are never edited in place.
|
||||
- In JSON mode the history is viewable in any text editor, and a manual rollback to a prior state is possible by removing final lines (one per transaction).
|
||||
|
||||
2. Differential persistence
|
||||
- Kanta stores diffs (patches), not full state, for normal writes.
|
||||
@@ -28,8 +28,7 @@ It is intentionally focused on the current standalone package behavior.
|
||||
|
||||
## On-Disk Record Types
|
||||
|
||||
Kanta uses a newline-delimited stream where each line is either a change
|
||||
record or a snapshot record.
|
||||
Kanta uses a newline-delimited stream where each line is either a change record or a snapshot record.
|
||||
|
||||
### Change record
|
||||
|
||||
@@ -68,29 +67,24 @@ Fields:
|
||||
3. Replay subsequent change records in order using patch application.
|
||||
4. The final replay state becomes in-memory `kanta.data`.
|
||||
|
||||
This model provides fast startup for large logs while retaining append-only
|
||||
history.
|
||||
This model provides fast startup for large logs while retaining append-only history.
|
||||
|
||||
## Serialization Semantics
|
||||
|
||||
- In-memory data is defined by an application `msgspec.Struct` type.
|
||||
- Kanta round-trips through plain builtins for persistence and diffing.
|
||||
- Dict keys are serialized as strings (`str_keys=True`) for stable JSON form.
|
||||
- Normalization changes introduced by struct decode/encode are logged together
|
||||
with migrations as `migrate:vN`, or as `migrate:msgspec` when no migration
|
||||
ran but normalization still produces a diff.
|
||||
- Normalization changes introduced by struct decode/encode are logged together with migrations as `migrate:vN`, or as `migrate:msgspec` when no migration ran but normalization still produces a diff.
|
||||
|
||||
## Transaction Semantics
|
||||
|
||||
- `kanta.transaction(action=...)` captures a pre-transaction snapshot dict.
|
||||
- By default a transaction updates the modification time `m` to the current UTC
|
||||
time.
|
||||
- By default a transaction updates the modification time `m` to the current UTC time.
|
||||
- `mtime=True|False|datetime` controls the modification time `m`:
|
||||
- `True` (default) sets `m` to the current UTC time.
|
||||
- `False` omits `m`, leaving the previous modification time in effect.
|
||||
- A `datetime` sets `m` to that explicit value.
|
||||
- System operations such as `migrate:msgspec` use `mtime=False` so they are not
|
||||
considered modifications and do not advance `m`.
|
||||
- System operations such as `migrate:msgspec` use `mtime=False` so they are not considered modifications and do not advance `m`.
|
||||
- On success:
|
||||
- compute diff between previous builtins and current builtins,
|
||||
- queue a `ChangeRecord` if non-empty,
|
||||
@@ -103,9 +97,7 @@ Nested transactions are rejected.
|
||||
|
||||
## Modification Time
|
||||
|
||||
`kanta.mtime` exposes the last modification time carried forward from change
|
||||
records. It is updated by normal transactions and preserved across snapshots and
|
||||
reloads, while system operations such as migrations leave it unchanged.
|
||||
`kanta.mtime` exposes the last modification time carried forward from change records. It is updated by normal transactions and preserved across snapshots and reloads, while system operations such as migrations leave it unchanged.
|
||||
|
||||
## Flush and Lifecycle
|
||||
|
||||
@@ -120,88 +112,52 @@ reloads, while system operations such as migrations leave it unchanged.
|
||||
- `await kanta.open()` (default) creates the database file if missing.
|
||||
- `await kanta.open(create=False)` fails when the file is missing or empty.
|
||||
- `await kanta.open(readonly=True)` opens an existing database read-only.
|
||||
- The file is opened without acquiring a lock and without a background flush
|
||||
task.
|
||||
- The file is opened without acquiring a lock and without a background flush task.
|
||||
- Existing records are replayed and migrations are still applied in memory.
|
||||
- Transactions and explicit flushes are rejected.
|
||||
- The file is never created if missing.
|
||||
|
||||
### Callbacks
|
||||
|
||||
All callbacks are registered via decorators and receive arguments by their
|
||||
annotation types. Parameters without a supported annotation are only allowed
|
||||
when they have a default value.
|
||||
All callbacks are registered via decorators and receive arguments by their annotation types. Parameters without a supported annotation are only allowed when they have a default value.
|
||||
|
||||
#### Bootstrap Callbacks
|
||||
|
||||
- When `open()` creates a new database, it always writes a single bootstrap
|
||||
`ChangeRecord`.
|
||||
- The simplest bootstrap is the initial data object passed to `Kanta(...)`;
|
||||
bootstrap callbacks are optional and only needed when you want to modify or
|
||||
enrich that object at creation time.
|
||||
- When `open()` creates a new database, it always writes a single bootstrap `ChangeRecord`.
|
||||
- The simplest bootstrap is the initial data object passed to `Kanta(...)`; bootstrap callbacks are optional and only needed when you want to modify or enrich that object at creation time.
|
||||
- Register callbacks via:
|
||||
- `@kanta.bootstrap`
|
||||
- `@kanta.bootstrap(action=..., user=..., mtime=...)`
|
||||
- Bootstrap callbacks may be sync or async. The live root data object is
|
||||
injected by annotating a parameter with the struct type passed to `Kanta`,
|
||||
and the `Kanta` instance itself can be injected by annotating a parameter
|
||||
with `Kanta`.
|
||||
- Bootstrap callbacks may be sync or async. The live root data object is injected by annotating a parameter with the struct type passed to `Kanta`, and the `Kanta` instance itself can be injected by annotating a parameter with `Kanta`.
|
||||
- Multiple bootstrap callbacks are supported:
|
||||
- callbacks execute in registration order,
|
||||
- exactly one bootstrap `ChangeRecord` is queued,
|
||||
- bootstrap metadata (`action`, `user`, `mtime`) is taken from the last
|
||||
callback registration.
|
||||
- If no bootstrap callbacks are registered, the bootstrap record still uses
|
||||
`action="bootstrap"` and contains the initial data object.
|
||||
- If any bootstrap callback raises, Kanta closes and removes the database file,
|
||||
then re-raises the exception.
|
||||
- bootstrap metadata (`action`, `user`, `mtime`) is taken from the last callback registration.
|
||||
- If no bootstrap callbacks are registered, the bootstrap record still uses `action="bootstrap"` and contains the initial data object.
|
||||
- If any bootstrap callback raises, Kanta closes and removes the database file, then re-raises the exception.
|
||||
|
||||
#### Fatal Error Handlers
|
||||
|
||||
- Fatal background persistence errors can be handled with `@kanta.fatal_error`.
|
||||
- Handlers may be sync or async. The `DatabaseError` is injected by annotating
|
||||
a parameter with `DatabaseError`; `Kanta` may also be injected.
|
||||
- Multiple handlers are supported and invoked in registration order. A failing
|
||||
handler is logged and does not prevent subsequent handlers from running.
|
||||
- Handlers may be sync or async. The `DatabaseError` is injected by annotating a parameter with `DatabaseError`; `Kanta` may also be injected.
|
||||
- Multiple handlers are supported and invoked in registration order. A failing handler is logged and does not prevent subsequent handlers from running.
|
||||
|
||||
#### Clock
|
||||
|
||||
- `@kanta.clock` registers a callback `() -> datetime` that replaces the
|
||||
default UTC clock. Its value is used for all record timestamps (`ts`, and
|
||||
`m` when `mtime` is `True`) and for snapshot timestamps.
|
||||
- The clock is only read when a timestamp is actually produced; no-op
|
||||
transactions and skipped snapshot checks do not read it.
|
||||
- Register before `open()` so that bootstrap and migration records use the
|
||||
custom clock as well. This is mainly useful for tests and reproducible
|
||||
demos.
|
||||
- `@kanta.clock` registers a callback `() -> datetime` that replaces the default UTC clock. Its value is used for all record timestamps (`ts`, and `m` when `mtime` is `True`) and for snapshot timestamps.
|
||||
- The clock is only read when a timestamp is actually produced; no-op transactions and skipped snapshot checks do not read it.
|
||||
- Register before `open()` so that bootstrap and migration records use the custom clock as well. This is mainly useful for tests and reproducible demos.
|
||||
|
||||
#### Transaction Log Formatting
|
||||
|
||||
- Logfmt callbacks prettify identifiers in the change log and are registered with
|
||||
`@kanta.logfmt`.
|
||||
- A logfmt callback is called for every value Kanta renders: diff values, path
|
||||
components, and the transaction `user`. It receives the value as its first
|
||||
parameter and optionally a `path: str` parameter with the dot-notation path
|
||||
to the value. The special path `"$user"` is used when rendering the
|
||||
transaction actor, replacing the old `user_display` parameter.
|
||||
- The callback returns `str | None`: a string replaces the default rendering,
|
||||
while `None` means "fall through to the next formatter".
|
||||
- State dicts are injected by parameter name or annotation tag, which share
|
||||
the same vocabulary: `prev` receives the previous state dict and `state`
|
||||
the current one. Matching by name ignores the annotation entirely. The
|
||||
`DictPrev`/`DictState` aliases (`Annotated[dict, "prev"]` /
|
||||
`Annotated[dict, "state"]`) work under any parameter name, and a tag takes
|
||||
precedence over the name. `DictPre` and `DictPost` are kept as aliases of
|
||||
`DictPrev` and `DictState`. The `Kanta` instance can also be injected.
|
||||
- Alternatively, a logfmt callback can be a class inheriting from `LogFmt`; the
|
||||
framework instantiates it with the state dicts and calls its
|
||||
`resolve(value, path) -> str | None` method.
|
||||
- Multiple logfmt callbacks are stacked in registration order; the first
|
||||
callback to return a non-`None` result wins. If none handle a value, Kanta
|
||||
falls back to its default formatting.
|
||||
- Logfmt callbacks prettify identifiers in the change log and are registered with `@kanta.logfmt`.
|
||||
- A logfmt callback is called for every value Kanta renders: diff values, path components, and the transaction `user`. It receives the value as its first parameter and optionally a `path: str` parameter with the dot-notation path to the value. The special path `"$user"` is used when rendering the transaction actor.
|
||||
- The callback returns `str | None`: a string replaces the default rendering, while `None` means "fall through to the next formatter".
|
||||
- State dicts are injected by parameter name or annotation tag, which share the same vocabulary: `prev` receives the previous state dict and `state` the current one. Matching by name ignores the annotation entirely. The `DictPrev`/`DictState` aliases (`Annotated[dict, "prev"]` / `Annotated[dict, "state"]`) work under any parameter name, and a tag takes precedence over the name. `DictPre` and `DictPost` are kept as aliases of `DictPrev` and `DictState`. The `Kanta` instance can also be injected.
|
||||
- Alternatively, a logfmt callback can be a class inheriting from `LogFmt`; the framework instantiates it with the state dicts and calls its `resolve(value, path) -> str | None` method.
|
||||
- Multiple logfmt callbacks are stacked in registration order; the first callback to return a non-`None` result wins. If none handle a value, Kanta falls back to its default formatting.
|
||||
|
||||
The decorator accepts an optional ``path`` so the callback only runs for
|
||||
values at that exact path:
|
||||
The decorator accepts an optional ``path`` so the callback only runs for values at that exact path:
|
||||
|
||||
```python
|
||||
@kanta.logfmt(path="$user")
|
||||
@@ -215,63 +171,21 @@ def resolve_user_key(value: str) -> str | None:
|
||||
|
||||
#### Transaction Log Headers
|
||||
|
||||
- By default a transaction is logged with an `action by user` header followed
|
||||
by the diff lines. Added paths are colored green, deleted paths red.
|
||||
- `kanta.transaction(..., extra=...)` accepts a display-only value that is
|
||||
shown after the action in the header. Anything other than `None` is
|
||||
printed str-converted (colored by Kanta), unless a custom logemit handler
|
||||
does something else with it; it is never persisted in the `ChangeRecord`.
|
||||
- `kanta.transaction(..., logdiff=False)` skips building and printing the diff
|
||||
body and logs only the header, which is useful for large or noisy
|
||||
changesets. Diff output can also be disabled globally with
|
||||
`configure_logging(diff=False)`; diff lines are emitted on the
|
||||
`kanta.transaction.diff` child logger so applications can route or silence
|
||||
them separately from the headers.
|
||||
- By default a transaction is logged with an `action by user` header followed by the diff lines. Added paths are colored green, deleted paths red.
|
||||
- `kanta.transaction(..., extra=...)` accepts a display-only value that is shown after the action in the header. Anything other than `None` is printed str-converted (colored by Kanta), unless a custom logemit handler does something else with it; it is never persisted in the `ChangeRecord`.
|
||||
- `kanta.transaction(..., logdiff=False)` skips building and printing the diff body and logs only the header, which is useful for large or noisy changesets. Diff output can also be disabled globally with `configure_logging(diff=False)`; diff lines are emitted on the `kanta.transaction.diff` child logger so applications can route or silence them separately from the headers.
|
||||
|
||||
#### Log Emitters
|
||||
|
||||
- Every change-related message Kanta emits (transaction/bootstrap/migration
|
||||
changes, file created/opened lines, migration summaries, aborted
|
||||
transactions) is described by a `kanta.logging.LogEvent` and dispatched
|
||||
through
|
||||
`kanta.logging.emit_event`. Kanta's own output goes through the same
|
||||
mechanism: when no `logemit` callback handles an event,
|
||||
`kanta.logging.default_emit` renders it with the built-in formatting.
|
||||
- A `LogEvent` carries the event `kind` (`"change"`, `"created"`,
|
||||
`"opened"`, `"migrated"`, `"aborted"`), the preferred `logger` and `level`,
|
||||
the
|
||||
`kanta` instance, and all relevant state: `action`, `user`, `extra`,
|
||||
`error` (for aborted transactions), `diff`, `previous`/`current` state
|
||||
dicts, the built `logfmt` chain, and version info for migration events.
|
||||
Application-specific context (e.g. a connection id) can be stored in
|
||||
`kanta.ctx` — a user-writable namespace — and read back in callbacks as
|
||||
`event.kanta.ctx`, which also covers creation/bootstrap events.
|
||||
- The built-in formatting is assembled from standard blocks that custom
|
||||
emitters can reuse as-is or replace piecemeal:
|
||||
- `event.header` — a lazy property producing the default one-line header
|
||||
for any kind: `<action>[ <extra>][ by <user>]` for changes,
|
||||
`<action>[ <extra>][ by <user>] transaction aborted: <error>` for aborts,
|
||||
and the `🛢️ <file> created|opened|migrated ...` summaries. It is
|
||||
settable: assign
|
||||
`event.header = ...` and return truthy to restyle the header while
|
||||
keeping the default diff routing.
|
||||
- `event.diff_lines` — a lazy property producing the pretty diff body for
|
||||
change events (built only if accessed).
|
||||
Every change-related message Kanta emits (transaction/bootstrap/migration changes, file created/opened lines, migration summaries, aborted transactions) is described by a `kanta.logging.LogEvent` and dispatched through `kanta.logging.emit_event`. Kanta's own output goes through the same mechanism: when no `logemit` callback handles an event, `kanta.logging.default_emit` renders it with the built-in formatting.
|
||||
|
||||
- A `LogEvent` carries the event `kind` (`"change"`, `"created"`, `"opened"`, `"migrated"`, `"aborted"`), the preferred `logger` and `level`, the `kanta` instance, and all relevant state: `action`, `user`, `extra`, `error` (for aborted transactions), `diff`, `previous`/`current` state dicts, the built `logfmt` chain, and version info for migration events. Application-specific context (e.g. a connection id) can be stored in `kanta.ctx` — a user-writable namespace — and read back in callbacks as `event.kanta.ctx`, which also covers creation/bootstrap events.
|
||||
- The built-in formatting is assembled from standard blocks that custom emitters can reuse as-is or replace piecemeal:
|
||||
- `event.header` — a lazy property producing the default one-line header for any kind: `<action>[ <extra>][ by <user>]` for changes, `<action>[ <extra>][ by <user>] transaction aborted: <error>` for aborts, and the `🛢️ <file> created|opened|migrated ...` summaries. It is settable: assign `event.header = ...` and return truthy to restyle the header while keeping the default diff routing.
|
||||
- `event.diff_lines` — a lazy property producing the pretty diff body for change events (built only if accessed).
|
||||
- `default_emit` itself is just `header` plus the `diff_lines` routing.
|
||||
- `@kanta.logemit` registers a callback receiving the event. The callback
|
||||
decides what is logged and where: it may log one or more messages on
|
||||
`event.logger`, log somewhere else, or nothing at all. A falsy return
|
||||
value marks the event handled and stops the chain; a truthy return value
|
||||
passes the event — possibly modified — to the next registered callback.
|
||||
When all callbacks pass, `default_emit` renders the event; a callback may
|
||||
also call `default_emit(event)` itself to delegate events it does not
|
||||
customize. Operational diagnostics (integrity errors, background flush
|
||||
failures) do not go through this mechanism.
|
||||
- Logging never breaks functionality: a crashing `logemit` callback is
|
||||
reported with `logger.exception` and the event falls back to the built-in
|
||||
formatting; if the built-in formatting itself fails, the error is reported
|
||||
and swallowed. The same applies to `logfmt` callbacks (a failing one is
|
||||
treated as a fall-through) and `logmigr` callbacks.
|
||||
- `@kanta.logemit` registers a callback receiving the event. The callback decides what is logged and where: it may log one or more messages on `event.logger`, log somewhere else, or nothing at all. A falsy return value marks the event handled and stops the chain; a truthy return value passes the event — possibly modified — to the next registered callback. When all callbacks pass, `default_emit` renders the event; a callback may also call `default_emit(event)` itself to delegate events it does not customize. Operational diagnostics (integrity errors, background flush failures) do not go through this mechanism.
|
||||
- Logging never breaks functionality: a crashing `logemit` callback is reported with `logger.exception` and the event falls back to the built-in formatting; if the built-in formatting itself fails, the error is reported and swallowed. The same applies to `logfmt` callbacks (a failing one is treated as a fall-through) and `logmigr` callbacks.
|
||||
|
||||
```python
|
||||
@kanta.logemit
|
||||
@@ -286,17 +200,9 @@ def emit(ev: LogEvent):
|
||||
#### Terminal Formatting Helpers
|
||||
|
||||
- `kanta.tty` provides the building blocks used by Kanta's own rendering:
|
||||
- `colors`: the mutable color palette. Colors are bare SGR parameter
|
||||
strings (e.g. `"1;34"`, `"38;5;226"`) without escape framing. Attributes
|
||||
are read at render time, so assignments (`colors.action = "36"`) and
|
||||
additions (`colors.session = "38;5;226"`) take effect immediately.
|
||||
- `Line`: builds a terminal string part by part. Calling it appends
|
||||
content (`str`-converted); `.<colorname>` arms a palette color for the
|
||||
next call only, and the reset is folded into a single escape sequence
|
||||
with whatever color comes next. `width=`/`align=` pad by display width;
|
||||
`str(line)` finishes the line and restores default colors.
|
||||
- `strip_ansi`, `displaywidth` (wide chars and emoji count correctly) and
|
||||
`pad` for working with pre-colored strings.
|
||||
- `colors`: the mutable color palette. Colors are bare SGR parameter strings (e.g. `"1;34"`, `"38;5;226"`) without escape framing. Attributes are read at render time, so assignments (`colors.action = "36"`) and additions (`colors.session = "38;5;226"`) take effect immediately.
|
||||
- `Line`: builds a terminal string part by part. Calling it appends content (`str`-converted); `.<colorname>` arms a palette color for the next call only, and the reset is folded into a single escape sequence with whatever color comes next. `width=`/`align=` pad by display width; `str(line)` finishes the line and restores default colors.
|
||||
- `strip_ansi`, `displaywidth` (wide chars and emoji count correctly) and `pad` for working with pre-colored strings.
|
||||
|
||||
## Migrations
|
||||
|
||||
@@ -308,8 +214,6 @@ def emit(ev: LogEvent):
|
||||
|
||||
## Safety Invariants
|
||||
|
||||
- Any detected out-of-transaction mutation is treated as a fatal consistency
|
||||
violation.
|
||||
- Any detected out-of-transaction mutation is treated as a fatal consistency violation.
|
||||
- Flush failures mark the instance as failed and trigger shutdown behavior.
|
||||
- Object identity of `kanta.data` is preserved across rollback when possible,
|
||||
minimizing stale-reference hazards for callers.
|
||||
- Object identity of `kanta.data` is preserved across rollback when possible, minimizing stale-reference hazards for callers.
|
||||
|
||||
@@ -0,0 +1,14 @@
|
||||
# Fatal Error Handlers
|
||||
|
||||
Fatal background write errors can be observed with a decorator:
|
||||
|
||||
```python
|
||||
import os
|
||||
import signal
|
||||
|
||||
@kanta.fatal_error
|
||||
async def on_fatal(err):
|
||||
os.kill(os.getpid(), signal.SIGTERM) # Die
|
||||
```
|
||||
|
||||
Multiple fatal handlers are supported and run in registration order.
|
||||
@@ -0,0 +1,12 @@
|
||||
# Migrations
|
||||
|
||||
Adding or removing a field and other such simple operations are automatic, but when the time comes to really change your data model, implement a `migrate_v1` function that converts your old data to the new form. This works on plain built-in dict and other types, to avoid needing to preserve old versions of your structs.
|
||||
|
||||
Pass a module (or import path) containing `migrate_vN` functions:
|
||||
|
||||
```python
|
||||
kanta = Kanta("data.kantadb", Data(), migrations="myapp.migrations")
|
||||
await kanta.open()
|
||||
```
|
||||
|
||||
Kanta tracks migration version metadata automatically, and fast forwards your database to current version by running all the migrations needed while opening the database.
|
||||
@@ -0,0 +1,70 @@
|
||||
# Database Rotation
|
||||
|
||||
Kanta can bound the on-disk history of a database to a configurable retention window (e.g. the last 30 days) by *rotating* the database file: the aged-out content is copied to a timestamped sibling file and the main file is truncated and rewritten in place with only the retained history plus fresh snapshots. Normal operation stays append-only under the exclusive lock; rotation is the only operation that rewrites the file.
|
||||
|
||||
## Configuration
|
||||
|
||||
Rotation is enabled with a keyword option on `Kanta(...)`:
|
||||
|
||||
- `retention: timedelta | int | None = None` — history window to keep; a plain `int` is interpreted as a number of days. `None` (default) disables rotation entirely.
|
||||
|
||||
Rotation uses the same clock as record timestamps, so a custom `@kanta.clock` callback controls it as well (useful in tests).
|
||||
|
||||
## When rotation runs
|
||||
|
||||
Rotation happens inside `Kanta.open()`, after acquiring the exclusive lock and before replay. At that point the file is quiescent: no background flush loop is running yet and no records are in flight. Rotation therefore never races with the background writer, snapshot requests, or migration snapshot writes.
|
||||
|
||||
Consequence: a database that is never reopened never rotates. For long-running services, rotation takes effect on the next restart.
|
||||
|
||||
## Rotated file naming
|
||||
|
||||
The history that aged out is preserved at:
|
||||
|
||||
```
|
||||
{stem}@{ISO-8601 timestamp}.kantadb
|
||||
```
|
||||
|
||||
- `{stem}` is the original filename with its extension stripped (`Path(filename).stem`), so databases named `data`, `data.kantadb`, or `data.db` all rotate to `data@….kantadb`.
|
||||
- The timestamp is the `ts` of the last record dropped by the rotation (the leading snapshot of the rewritten main file carries the same `ts`), not the rotation time — the name tells you exactly which point in history the rotated file ends at. It is rendered in ISO 8601 basic format at second precision (e.g. `20260902T143000Z`); the exact microsecond timestamp of the cutoff remains available inside the file as the `ts` of its final line and of the leading snapshot of the new file.
|
||||
- Rotated files live in the same directory as the main file.
|
||||
- On collision (a rotated file with the same name already exists), an incrementing suffix is inserted before the extension (`data@….1.kantadb`, `data@….2.kantadb`, …) rather than overwriting.
|
||||
- Rotated files are never deleted by rotation.
|
||||
|
||||
## Rotation algorithm
|
||||
|
||||
Let `cutoff = now - retention`. All planning happens on the in-memory bytes of the database file already read by open; no second disk read is needed.
|
||||
|
||||
1. **Eligibility.** Rotation is skipped when there is nothing to do: when the file contains no change records older than `cutoff` (the retention window already covers all history), or when the file contains no change records at all (a snapshot-only file is treated as already fully rotated and never re-rotated).
|
||||
|
||||
2. **Replay base.** The base is the newest snapshot whose `ts <= cutoff`, falling back to the start of file when no such snapshot exists. Starting at the most recent snapshot would silently drop history that must be retained.
|
||||
|
||||
3. **Replay and validate.** The file is replayed from the base forward to the end. Records with `ts < cutoff` are applied to the replay (they are needed to reach the cutoff state) but not retained in the output. At every snapshot encountered after the base, the replayed state is validated against the snapshot state; a mismatch means the history is corrupt, and rotation is aborted with a `DatabaseError`, leaving the original file untouched. The byte offset just after the last record with `ts < cutoff` (frame-boundary aligned) is remembered as `cutoff_end`.
|
||||
|
||||
4. **Copy the original aside.** The main file is copied with `shutil.copy2` directly to its final `{stem}@{ts}.kantadb` name. On filesystems with copy-on-write this performs a cheap reflink copy. A failure here aborts rotation with the original file intact.
|
||||
|
||||
5. **Rewrite the main file in place.** On the locked file descriptor the content is replaced (seek 0, truncate, write, fsync) with, in order:
|
||||
1. A **snapshot of the state at the cutoff**, stamped with the schema version and modification time in effect at the cutoff. Its `ts` is the `ts` of the last pre-cutoff record — the same timestamp used in the rotated filename. This snapshot is the new replay base and carries the version forward so migrations are not re-run; it is always written.
|
||||
2. The retained change records (`ts >= cutoff`), recreated record by record — no internal snapshots are carried over.
|
||||
3. A **final snapshot** of the state after the last retained record, written only when enough changes were retained to warrant one (the same policy as regular snapshot writes). If no records survived the cutoff, the new file consists of the single leading snapshot and nothing else — the steady state for databases whose history has fully aged out.
|
||||
|
||||
6. **Trim the rotated copy.** The rotated file is truncated to `cutoff_end` bytes, so it contains only the dropped history and does not duplicate the records retained in the main file. The cut is at a frame boundary, so the rotated file remains a valid, replayable database on its own. This happens only after step 5's fsync, so until then the rotated file still holds the complete original content as a crash-recovery anchor.
|
||||
|
||||
7. **Continue normal open.** Replay and migrations proceed on the same locked file. Because the leading snapshot carries the current version, migrations run exactly as they would have against the old content.
|
||||
|
||||
Failure rule: any error before step 5 leaves the main file byte-identical (only an extra copy exists). A crash during step 5 may leave the main file torn, but the rotated copy still holds the complete original content — recovery is copying it back. After step 6 the split is complete and both files are consistent.
|
||||
|
||||
## Design notes
|
||||
|
||||
**In-place rewrite, not rename-and-recreate.** The writer holds an exclusive `flock` on the file from `open()` until `close()`, and `flock` is attached to the open file description (inode), not the path. Renaming the locked file away and creating a fresh file at the main path would open a race: between the rename and the creation, a second instance could open the missing path with `O_CREAT`, acquire its own lock on the fresh inode, and bootstrap a divergent database. On Windows, renaming the locked file would fail outright (the database is opened with `FILE_SHARE_READ` only). Rotation therefore never renames or unlinks the main file and never releases its lock; a second instance opening the path at any moment gets either the old content or the new, and never its own lock. The only primitive this requires is `LockedFile.replace_content()` (seek 0, truncate, write, fsync).
|
||||
|
||||
**Rewrite (re-frame), not verbatim copy**, for all records written to the main file. `BinFramer` checksums are offset-keyed (the checksum includes the absolute `record_offset`), so a verbatim byte copy to a new offset would be unreadable; binary records are re-framed at their new offsets. Rewriting also normalizes encoding drift and lets rotation drop the redundant intermediate snapshots the original file accumulated. The rotated copy is the one place where verbatim bytes are used — a raw `copy2` plus a frame-aligned tail truncation — which is safe precisely because it preserves original offsets: the truncated prefix keeps every frame at its original `record_offset`, so binary checksums stay valid.
|
||||
|
||||
## Integrity guarantees
|
||||
|
||||
- Rotation runs under the exclusive lock, before the background writer starts.
|
||||
- The main path is never renamed, unlinked, or unlocked during rotation; no bootstrap race with a second instance is possible.
|
||||
- Replay is validated against every snapshot in range; any validation failure aborts rotation with the original file intact.
|
||||
- A leading cutoff snapshot (ts = last pre-cutoff record, schema version at the cutoff) is always written; a final snapshot is written only when warranted by retained changes.
|
||||
- The full original content sits at `{stem}@{ts}.kantadb` before the main file is touched, and is only trimmed to the dropped-history prefix after the rewritten main file is fsynced.
|
||||
- Rotated files are never deleted by rotation.
|
||||
- Files with no change records (already reduced to a snapshot) are never re-rotated.
|
||||
@@ -0,0 +1,96 @@
|
||||
# Usage Patterns: Opening and Data Handling
|
||||
|
||||
This document covers the basics of opening a database and working with the live data object: who owns it, where initial data comes from, and the lifecycle patterns Kanta supports.
|
||||
|
||||
## The data object belongs to you
|
||||
|
||||
You pass the root state object to `Kanta(...)`, and Kanta never replaces it with a new instance. Replay, migrations, transaction rollbacks — all restore or mutate the object's contents in place, preserving its identity. This means you may hold external references to the object (or to parts of it) and they stay valid for the lifetime of the `Kanta` instance:
|
||||
|
||||
```python
|
||||
data = Data()
|
||||
kanta = Kanta("data.kantadb", data)
|
||||
await kanta.open()
|
||||
|
||||
assert kanta.data is data # always the same object
|
||||
```
|
||||
|
||||
Access the state as `kanta.data`, via your own reference, or both — they are the same object. Mutations must happen inside a transaction (see below); Kanta treats any detected out-of-transaction mutation as a fatal consistency violation.
|
||||
|
||||
Note that replacing the whole object is also possible (`kanta.data = Data()` has a setter), but then previously held references point at the old object — prefer in-place mutation.
|
||||
|
||||
## Where initial data comes from
|
||||
|
||||
The object passed to `Kanta(...)` seeds the database: when `open()` creates a brand-new (missing or empty) file, it writes a single bootstrap change record containing that object's contents, optionally modified by `@kanta.bootstrap` handlers (see [Bootstrap and open modes](bootstrap.md)).
|
||||
|
||||
When the file already exists, the constructor argument is *not* used as state — replay rebuilds the contents of `self.data` in place from the stored records, and the argument only defines the struct type. Consequence: if you close an instance, delete the file, and `open()` again, the new database is bootstrapped from the object's *current* contents — the old database's final state, not its original initial values. For a genuinely fresh start, construct a new data object (or reset the fields in a bootstrap handler).
|
||||
|
||||
## Lifecycle patterns
|
||||
|
||||
### Context manager (single database, scoped lifetime)
|
||||
|
||||
```python
|
||||
async with Kanta("data.kantadb", Data()) as kanta:
|
||||
with kanta.transaction(action="create_user") as data:
|
||||
data.users[user_id] = User(name="Alice")
|
||||
```
|
||||
|
||||
`async with` guarantees `open()`/`close()` pairing: the final flush happens on exit even on errors. Both `as` bindings are pure shortcuts — `async with Kanta(...) as kanta` binds the `Kanta` instance itself, and `with kanta.transaction(...) as data` binds exactly `kanta.data`. Use them or don't:
|
||||
|
||||
```python
|
||||
async with Kanta("data.kantadb", Data()) as db:
|
||||
with db.transaction(action="rename"):
|
||||
db.data.users[user_id].name = "Bob" # same object as `data` above
|
||||
```
|
||||
|
||||
### Module-level global (typical application state)
|
||||
|
||||
When the database lives as long as the process, define it once and open/close at application startup and shutdown:
|
||||
|
||||
```python
|
||||
kanta = Kanta("data.kantadb", Data())
|
||||
data = kanta.data # optional: your own direct reference
|
||||
|
||||
|
||||
async def startup() -> None:
|
||||
await kanta.open()
|
||||
|
||||
|
||||
async def shutdown() -> None:
|
||||
await kanta.close()
|
||||
|
||||
|
||||
async def create_user(name: str) -> None:
|
||||
with kanta.transaction(action="create_user") as d:
|
||||
d.users[uuid7()] = User(name=name)
|
||||
```
|
||||
|
||||
Transactions are synchronous context managers, so no `await` is needed per operation; the background task flushes queued changes periodically.
|
||||
|
||||
### Dynamically created databases (per-project, per-tenant, …)
|
||||
|
||||
Nothing requires module-level definitions. An app managing many databases simply constructs instances on demand and tracks them itself:
|
||||
|
||||
```python
|
||||
class ProjectStore:
|
||||
def __init__(self) -> None:
|
||||
self.projects: dict[str, Kanta[Data]] = {}
|
||||
|
||||
async def get(self, project_id: str) -> Kanta[Data]:
|
||||
kanta = self.projects.get(project_id)
|
||||
if kanta is None:
|
||||
kanta = Kanta(f"projects/{project_id}.kantadb", Data())
|
||||
await kanta.open()
|
||||
self.projects[project_id] = kanta
|
||||
return kanta
|
||||
```
|
||||
|
||||
Remember that each open writer holds an exclusive lock on its file, so keep one `Kanta` instance per path and close instances you no longer need.
|
||||
|
||||
## Transactions in brief
|
||||
|
||||
- `with kanta.transaction(action=..., user=...) as data:` yields the live state object for mutation.
|
||||
- On success, Kanta computes a diff against the pre-transaction state and queues a change record; on exception, the in-memory data is rolled back and the exception is re-raised.
|
||||
- Nested transactions are rejected.
|
||||
- A transaction that changes nothing queues no record.
|
||||
|
||||
See [On-disk format](database.md) for the full transaction semantics and [Validation](validation.md) for integrity checks that run inside each transaction.
|
||||
@@ -0,0 +1,20 @@
|
||||
# Validation
|
||||
|
||||
`@kanta.validate` registers an integrity validator for the database state, beyond the structural validation msgspec already performs during decoding.
|
||||
|
||||
```python
|
||||
@kanta.validate
|
||||
def check_users(data: Data) -> None:
|
||||
for user in data.users.values():
|
||||
if not user.name:
|
||||
raise ValueError("user without a name")
|
||||
```
|
||||
|
||||
Validators receive the live data object (and optionally the `Kanta` instance as a second annotated parameter) and must **raise an exception** when the data is inconsistent. They are not intended to modify or correct the data — only to fail.
|
||||
|
||||
Validators run:
|
||||
|
||||
- **on open** — after replay, msgspec decoding and migrations, before the database becomes usable; a failure aborts the open and releases the file,
|
||||
- **after each transaction** — before the change is committed to history; a failure rolls the transaction back, so invalid state never reaches the log.
|
||||
|
||||
Multiple validators may be registered; they run in registration order until the first failure. Validators must be synchronous (transactions are synchronous) — async callbacks are rejected at registration time.
|
||||
+19
-3
@@ -136,6 +136,7 @@ class CallbackRegistry:
|
||||
"bootstrap": [],
|
||||
"fatal_error": [],
|
||||
"logmigr": [],
|
||||
"validate": [],
|
||||
}
|
||||
self._logfmt_callbacks: list[_LogFmtFunctionSpec | _LogFmtClassSpec] = []
|
||||
self._logemit_callbacks: list[Callable[..., Any]] = []
|
||||
@@ -174,6 +175,10 @@ class CallbackRegistry:
|
||||
raise TypeError(f"{kind} callbacks must be functions, not classes")
|
||||
if not callable(callback):
|
||||
raise TypeError(f"{kind} callback must be callable")
|
||||
if kind == "validate" and inspect.iscoroutinefunction(callback):
|
||||
raise TypeError(
|
||||
"validate callbacks must not be async (transactions are synchronous)"
|
||||
)
|
||||
|
||||
params = self._validate_function(callback, kind)
|
||||
is_async = inspect.iscoroutinefunction(callback)
|
||||
@@ -215,6 +220,16 @@ class CallbackRegistry:
|
||||
break
|
||||
return results
|
||||
|
||||
def invoke_sync(self, kind: str, ctx: InjectionContext) -> None:
|
||||
"""Invoke all sync callbacks of *kind* in order; first exception raises.
|
||||
|
||||
Used for ``validate`` callbacks, which run inside synchronous
|
||||
transactions and therefore must not be async.
|
||||
"""
|
||||
for reg in self._callbacks[kind]:
|
||||
kwargs = self._build_kwargs(reg.params, ctx)
|
||||
reg.callback(**kwargs)
|
||||
|
||||
def has(self, kind: str) -> bool:
|
||||
"""Return True if any callback of *kind* is registered."""
|
||||
if kind == "logfmt":
|
||||
@@ -531,22 +546,23 @@ class CallbackRegistry:
|
||||
if bare is MigrationReport:
|
||||
return kind == "logmigr"
|
||||
if self._data_type is not None and bare is self._data_type:
|
||||
return kind == "bootstrap"
|
||||
return kind in {"bootstrap", "validate"}
|
||||
if self._kanta_class is not None and bare is self._kanta_class:
|
||||
return kind in {
|
||||
"bootstrap",
|
||||
"fatal_error",
|
||||
"logfmt",
|
||||
"logmigr",
|
||||
"validate",
|
||||
}
|
||||
return False
|
||||
|
||||
def _allowed_message(self, kind: str) -> str:
|
||||
parts: list[str] = []
|
||||
if kind == "bootstrap":
|
||||
if kind in {"bootstrap", "validate"}:
|
||||
if self._data_type is not None:
|
||||
parts.append(self._data_type.__name__)
|
||||
if kind in {"bootstrap", "fatal_error", "logfmt", "logmigr"}:
|
||||
if kind in {"bootstrap", "fatal_error", "logfmt", "logmigr", "validate"}:
|
||||
if self._kanta_class is not None:
|
||||
parts.append(self._kanta_class.__name__)
|
||||
if kind == "fatal_error":
|
||||
|
||||
+49
-46
@@ -1,63 +1,66 @@
|
||||
"""Diff computation and replay utilities."""
|
||||
"""Diff computation and replay utilities.
|
||||
|
||||
import jsondiff
|
||||
Diff format: JSON-serializable dicts where ``$delete`` is the only command
|
||||
our producer emits; added keys and changed values (scalars, lists, type
|
||||
changes — lists always wholesale) are plain assignments. A dict value
|
||||
assigned over a non-dict needs no ``$replace``: the consumer can see from
|
||||
the old value whether to patch (old is a dict) or replace. User keys
|
||||
starting with ``$`` are escaped by prepending another ``$``
|
||||
(``$foo`` -> ``$$foo``); values are stored verbatim.
|
||||
|
||||
The consumer additionally stays compatible with jsondiff's marshaled
|
||||
syntax, so it can replay diffs produced by jsondiff itself: ``$replace``,
|
||||
positional ``$insert``/``$delete`` and per-index nested diffs on lists,
|
||||
and jsondiff's escaping of ``$``-prefixed values.
|
||||
"""
|
||||
|
||||
from kanta.structs import ChangeRecord
|
||||
from kanta.serialization.base import ReplayResult, replay
|
||||
from kanta.serialization.base import ReplayResult, apply_diff, replay
|
||||
from kanta.serialization.framing import LineFramer
|
||||
from kanta.serialization.json import JsonSerializer
|
||||
|
||||
_UNCHANGED = object()
|
||||
|
||||
|
||||
def _escape_key(key: str) -> str:
|
||||
"""Escape a user key for use as a diff key (``$foo`` -> ``$$foo``)."""
|
||||
if isinstance(key, str) and key.startswith("$"):
|
||||
return "$" + key
|
||||
return key
|
||||
|
||||
|
||||
def _diff(previous, current):
|
||||
"""Compute a raw diff, or _UNCHANGED if there is no difference."""
|
||||
if isinstance(previous, dict) and isinstance(current, dict):
|
||||
result = {}
|
||||
deleted = [_escape_key(k) for k in previous if k not in current]
|
||||
if deleted:
|
||||
result["$delete"] = deleted
|
||||
for key, new_value in current.items():
|
||||
if key not in previous:
|
||||
result[_escape_key(key)] = new_value
|
||||
else:
|
||||
sub = _diff(previous[key], new_value)
|
||||
if sub is not _UNCHANGED:
|
||||
result[_escape_key(key)] = sub
|
||||
return result if result else _UNCHANGED
|
||||
if previous == current:
|
||||
return _UNCHANGED
|
||||
return current
|
||||
|
||||
|
||||
def compute_diff(previous: dict, current: dict) -> dict | None:
|
||||
"""Compute a jsondiff patch between two dicts.
|
||||
"""Compute a marshaled diff between two state dicts.
|
||||
|
||||
Returns None if there is no difference.
|
||||
"""
|
||||
return jsondiff.diff(previous, current, marshal=True) or None
|
||||
|
||||
|
||||
def _apply_diff(state: dict, diff: dict) -> dict:
|
||||
"""Apply a jsondiff patch manually, handling ``$replace`` and ``$delete``.
|
||||
|
||||
jsondiff.patch does not handle nested ``$replace`` commands when the
|
||||
parent key is missing from the state. This function recursively applies
|
||||
diffs, treating ``$replace`` as full replacement and ``$delete`` as
|
||||
key removal.
|
||||
"""
|
||||
if not isinstance(diff, dict):
|
||||
return diff
|
||||
|
||||
result = dict(state) if isinstance(state, dict) else state
|
||||
if not isinstance(result, dict):
|
||||
result = {}
|
||||
|
||||
for key, value in diff.items():
|
||||
if key == "$replace":
|
||||
return value
|
||||
elif key == "$delete":
|
||||
if isinstance(value, list):
|
||||
for k in value:
|
||||
result.pop(k, None)
|
||||
else:
|
||||
result.pop(value, None)
|
||||
elif isinstance(value, dict):
|
||||
old = result.get(key, {})
|
||||
if not isinstance(old, dict):
|
||||
old = {}
|
||||
result[key] = _apply_diff(old, value)
|
||||
else:
|
||||
result[key] = value
|
||||
|
||||
return result
|
||||
diff = _diff(previous, current)
|
||||
return diff if diff is not _UNCHANGED else None
|
||||
|
||||
|
||||
def patch_state(state: dict, diff: dict) -> dict:
|
||||
"""Apply a jsondiff patch to a state dict.
|
||||
|
||||
The diff was produced with ``marshal=True`` (string keys like
|
||||
``"$replace"`` and ``"$delete"``) and decoded from JSON.
|
||||
"""
|
||||
return _apply_diff(state, diff)
|
||||
"""Apply a marshaled diff to a state dict."""
|
||||
return apply_diff(state, diff)
|
||||
|
||||
|
||||
# Backward-compatible JSONL replay using the default serializer.
|
||||
|
||||
@@ -83,6 +83,10 @@ if sys.platform == "win32":
|
||||
]
|
||||
_kernel32.CloseHandle.restype = wintypes.BOOL
|
||||
_kernel32.CloseHandle.argtypes = [wintypes.HANDLE]
|
||||
_kernel32.SetEndOfFile.restype = wintypes.BOOL
|
||||
_kernel32.SetEndOfFile.argtypes = [wintypes.HANDLE]
|
||||
_kernel32.FlushFileBuffers.restype = wintypes.BOOL
|
||||
_kernel32.FlushFileBuffers.argtypes = [wintypes.HANDLE]
|
||||
|
||||
def _is_invalid_handle(handle) -> bool:
|
||||
return ctypes.c_void_p(handle).value == ctypes.c_void_p(-1).value
|
||||
@@ -177,6 +181,22 @@ class LockedFile:
|
||||
os.lseek(self._fd, current, os.SEEK_SET)
|
||||
return end
|
||||
|
||||
def replace_content(self, data: bytes) -> None:
|
||||
"""Atomically-ish rewrite the file's content in place, lock retained.
|
||||
|
||||
Seeks to the start, truncates, writes *data* and fsyncs, all on the
|
||||
already-locked descriptor. The path is never unlinked or renamed, so
|
||||
no other process can observe a missing file or acquire its own lock.
|
||||
Used by database rotation.
|
||||
"""
|
||||
if self._fd is None:
|
||||
raise RuntimeError("LockedFile.replace_content() called on a closed file")
|
||||
|
||||
if sys.platform == "win32":
|
||||
self._replace_content_win32(data)
|
||||
else:
|
||||
self._replace_content_unix(data)
|
||||
|
||||
def close(self) -> None:
|
||||
"""Release the lock and close the file."""
|
||||
if self._fd is None:
|
||||
@@ -227,6 +247,15 @@ class LockedFile:
|
||||
os.lseek(self._fd, 0, os.SEEK_END)
|
||||
os.write(self._fd, data)
|
||||
|
||||
def _replace_content_unix(self, data: bytes) -> None:
|
||||
os.lseek(self._fd, 0, os.SEEK_SET)
|
||||
os.ftruncate(self._fd, 0)
|
||||
view = memoryview(data)
|
||||
while view:
|
||||
written = os.write(self._fd, view)
|
||||
view = view[written:]
|
||||
os.fdatasync(self._fd)
|
||||
|
||||
# -- Windows -------------------------------------------------------------
|
||||
|
||||
def _open_win32(self, path: Path, create: bool, readonly: bool) -> None:
|
||||
@@ -288,3 +317,24 @@ class LockedFile:
|
||||
)
|
||||
if not ok:
|
||||
raise OSError(f"WriteFile failed: Windows error {ctypes.get_last_error()}")
|
||||
|
||||
def _replace_content_win32(self, data: bytes) -> None:
|
||||
_kernel32.SetFilePointer(self._fd, 0, None, _FILE_BEGIN)
|
||||
written = wintypes.DWORD()
|
||||
ok = _kernel32.WriteFile(
|
||||
self._fd,
|
||||
data,
|
||||
len(data),
|
||||
ctypes.byref(written),
|
||||
None,
|
||||
)
|
||||
if not ok:
|
||||
raise OSError(f"WriteFile failed: Windows error {ctypes.get_last_error()}")
|
||||
if not _kernel32.SetEndOfFile(self._fd):
|
||||
raise OSError(
|
||||
f"SetEndOfFile failed: Windows error {ctypes.get_last_error()}"
|
||||
)
|
||||
if not _kernel32.FlushFileBuffers(self._fd):
|
||||
raise OSError(
|
||||
f"FlushFileBuffers failed: Windows error {ctypes.get_last_error()}"
|
||||
)
|
||||
|
||||
+31
-1
@@ -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 | int | None = None,
|
||||
):
|
||||
"""Initialize a Kanta persistence instance.
|
||||
|
||||
@@ -63,6 +64,13 @@ 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, either a
|
||||
:class:`~datetime.timedelta` or a plain number of days. 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 +86,7 @@ class Kanta(Generic[T]):
|
||||
type=data_type,
|
||||
migrations=migrations,
|
||||
flush_interval=flush_interval,
|
||||
retention=retention,
|
||||
kanta=self,
|
||||
)
|
||||
|
||||
@@ -240,6 +249,27 @@ class Kanta(Generic[T]):
|
||||
return _register
|
||||
return _register(fn)
|
||||
|
||||
def validate(self, fn):
|
||||
"""Register a data validation callback.
|
||||
|
||||
Used as ``@kanta.validate``. The callback receives the live data
|
||||
object (and optionally the ``Kanta`` instance) and must raise an
|
||||
exception when the data is inconsistent. Validators run after replay
|
||||
during :meth:`open` (after msgspec decoding and migrations) and after
|
||||
each transaction, before the change is committed to history. Multiple
|
||||
validators run in registration order until the first failure.
|
||||
|
||||
Validators must be synchronous and must not modify the data — they
|
||||
only fail. A failure inside a transaction rolls the transaction back;
|
||||
a failure during open aborts the open.
|
||||
"""
|
||||
|
||||
def _register(callback):
|
||||
self._impl.add_validate(callback)
|
||||
return callback
|
||||
|
||||
return _register(fn)
|
||||
|
||||
def fatal_error(self, fn=None):
|
||||
"""Register fatal error handler callback.
|
||||
|
||||
|
||||
+53
-1
@@ -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,10 @@ class KantaImpl(PersistenceMixin, Generic[T]):
|
||||
self.data: T = kwargs.pop("data")
|
||||
self._kanta = kwargs.pop("kanta", None)
|
||||
migrations = kwargs.pop("migrations", None)
|
||||
retention = kwargs.pop("retention", None)
|
||||
if isinstance(retention, int) and not isinstance(retention, bool):
|
||||
retention = timedelta(days=retention)
|
||||
self.retention: timedelta | None = retention
|
||||
self.ctx = SimpleNamespace()
|
||||
super().__init__(**kwargs)
|
||||
self.migrations: Migrations | None = None
|
||||
@@ -91,6 +96,10 @@ class KantaImpl(PersistenceMixin, Generic[T]):
|
||||
"""Register one migration logging callback."""
|
||||
self.callback_registry.register("logmigr", callback)
|
||||
|
||||
def add_validate(self, callback) -> None:
|
||||
"""Register one data validation callback."""
|
||||
self.callback_registry.register("validate", callback)
|
||||
|
||||
def add_logemit(self, callback) -> None:
|
||||
"""Register one log emitter callback."""
|
||||
self.callback_registry.register("logemit", callback)
|
||||
@@ -181,6 +190,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(
|
||||
@@ -239,6 +251,16 @@ class KantaImpl(PersistenceMixin, Generic[T]):
|
||||
self.data_type,
|
||||
serializer=self.serializer,
|
||||
)
|
||||
if self.callback_registry.has("validate"):
|
||||
try:
|
||||
self.callback_registry.invoke_sync(
|
||||
"validate",
|
||||
InjectionContext(data=self.data, kanta=self._kanta),
|
||||
)
|
||||
except Exception:
|
||||
self.opened = False
|
||||
self.file.close()
|
||||
raise
|
||||
self.version = rr.version
|
||||
self.mtime = rr.m
|
||||
if log is not False and not migrations_ran:
|
||||
@@ -314,6 +336,12 @@ class KantaImpl(PersistenceMixin, Generic[T]):
|
||||
InjectionContext(data=self.data, kanta=self._kanta),
|
||||
)
|
||||
|
||||
if self.callback_registry.has("validate"):
|
||||
self.callback_registry.invoke_sync(
|
||||
"validate",
|
||||
InjectionContext(data=self.data, kanta=self._kanta),
|
||||
)
|
||||
|
||||
self.statedict = {}
|
||||
current = struct_to_dict(self.data, serializer=self.serializer)
|
||||
record = self.queue_change(
|
||||
@@ -375,6 +403,30 @@ 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:
|
||||
|
||||
+17
-2
@@ -15,6 +15,7 @@ from typing import Any
|
||||
|
||||
import msgspec
|
||||
|
||||
from kanta.serialization.base import _apply, unmarshal
|
||||
from kanta.tty import Line, displaywidth
|
||||
|
||||
transaction_logger = logging.getLogger("kanta.transaction")
|
||||
@@ -312,6 +313,13 @@ def _collect_changes(
|
||||
changes.append(("update" if existed else "add", path, diff))
|
||||
return
|
||||
|
||||
old_at_path = _get_nested(previous, path)
|
||||
if isinstance(old_at_path, list):
|
||||
# List edits ($insert/$delete/per-index) are shown as one whole-list
|
||||
# update; the diff is already unmarshaled at this point.
|
||||
changes.append(("update", path, _apply(old_at_path, diff)))
|
||||
return
|
||||
|
||||
for key, value in diff.items():
|
||||
if key == "$delete":
|
||||
if isinstance(value, list):
|
||||
@@ -340,7 +348,14 @@ def _collect_changes(
|
||||
("update" if old_collection is not None else "add", path, value)
|
||||
)
|
||||
elif isinstance(key, str) and key.startswith("$"):
|
||||
changes.append(("add", path, {key: value}))
|
||||
# Unknown $-command or (post-unmarshal) a user key starting with
|
||||
# "$": treat as a normal key.
|
||||
new_path = path + [str(key)]
|
||||
existed = _get_nested(previous, new_path) is not None
|
||||
if existed:
|
||||
_collect_changes(value, new_path, changes, previous)
|
||||
else:
|
||||
changes.append(("add", new_path, value))
|
||||
else:
|
||||
new_path = path + [str(key)]
|
||||
existed = _get_nested(previous, new_path) is not None
|
||||
@@ -418,7 +433,7 @@ def format_diff(
|
||||
Returns a list of formatted lines (without newlines).
|
||||
"""
|
||||
changes: list[tuple[str, list[str], Any]] = []
|
||||
_collect_changes(diff, [], changes, previous)
|
||||
_collect_changes(unmarshal(diff), [], changes, previous)
|
||||
if not changes:
|
||||
return []
|
||||
lines = []
|
||||
|
||||
@@ -0,0 +1,225 @@
|
||||
"""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 at second precision (e.g.
|
||||
``20260902T143000Z``); the exact microsecond timestamp remains available
|
||||
inside the file if ever needed. On collision an incrementing suffix is
|
||||
inserted before the extension.
|
||||
"""
|
||||
stamp = ts.strftime("%Y%m%dT%H%M%SZ")
|
||||
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
|
||||
|
||||
# The cutoff state starts from the replay base: when the base snapshot
|
||||
# already predates the cutoff, it may itself be the cutoff state.
|
||||
state_at_cutoff: dict[str, Any] | None = (
|
||||
copy.deepcopy(state) if base is not None else 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
|
||||
+87
-20
@@ -121,33 +121,100 @@ def replay(
|
||||
|
||||
|
||||
def _patch_state(state: dict, diff: dict) -> dict:
|
||||
return _apply_diff(state, diff)
|
||||
return apply_diff(state, diff)
|
||||
|
||||
|
||||
def _apply_diff(state: dict, diff: dict) -> dict:
|
||||
def _unescape(value: str) -> str:
|
||||
"""Reverse jsondiff's ``$$`` escaping; command strings pass through.
|
||||
|
||||
Only a ``$$`` prefix is stripped: jsondiff escapes ``$x`` to ``$$x``,
|
||||
while single ``$`` strings occur verbatim in our own diffs (we do not
|
||||
escape values) and must be left alone.
|
||||
"""
|
||||
if value.startswith("$$"):
|
||||
return value[1:]
|
||||
return value
|
||||
|
||||
|
||||
def unmarshal(diff: Any) -> Any:
|
||||
"""Unescape a marshaled diff (keys, values and ``$delete`` entries).
|
||||
|
||||
Needed for jsondiff-produced diffs, which escape ``$``-prefixed values
|
||||
as well as keys; our own producer escapes keys only, so unescaping
|
||||
values is a no-op for them.
|
||||
"""
|
||||
if isinstance(diff, dict):
|
||||
return {
|
||||
_unescape(k) if isinstance(k, str) else k: unmarshal(v)
|
||||
for k, v in diff.items()
|
||||
}
|
||||
if isinstance(diff, list):
|
||||
return [unmarshal(v) for v in diff]
|
||||
if isinstance(diff, str):
|
||||
return _unescape(diff)
|
||||
return diff
|
||||
|
||||
|
||||
def apply_diff(state: Any, diff: Any) -> Any:
|
||||
"""Apply a diff.
|
||||
|
||||
Understands our own format (plain assignment + ``$delete``) and
|
||||
jsondiff's marshaled syntax: ``$replace``, positional
|
||||
``$delete``/``$insert`` and per-index nested diffs on lists. A bare
|
||||
dict over a non-dict old value is a wholesale replacement.
|
||||
"""
|
||||
return _apply(state, unmarshal(diff))
|
||||
|
||||
|
||||
def _is_list_patch(diff: dict) -> bool:
|
||||
"""Whether a dict diff against a list state is a jsondiff list edit."""
|
||||
for key in diff:
|
||||
if key in ("$delete", "$insert"):
|
||||
continue
|
||||
try:
|
||||
int(key)
|
||||
except (ValueError, TypeError):
|
||||
return False
|
||||
return True
|
||||
|
||||
|
||||
def _apply(state: Any, diff: Any) -> Any:
|
||||
if not isinstance(diff, dict):
|
||||
return diff
|
||||
if not diff:
|
||||
return state
|
||||
if "$replace" in diff:
|
||||
return diff["$replace"]
|
||||
|
||||
result = dict(state) if isinstance(state, dict) else state
|
||||
if not isinstance(result, dict):
|
||||
result = {}
|
||||
if isinstance(state, list):
|
||||
if not _is_list_patch(diff):
|
||||
# Our own producer replaces a list with a dict (or any other
|
||||
# type) by plain assignment — no $replace wrapper.
|
||||
return diff
|
||||
result = list(state)
|
||||
deletes = diff.get("$delete")
|
||||
if deletes:
|
||||
for pos in deletes:
|
||||
result.pop(pos)
|
||||
for pos, value in diff.get("$insert", []):
|
||||
result.insert(pos, value)
|
||||
for key, value in diff.items():
|
||||
if key in ("$delete", "$insert"):
|
||||
continue
|
||||
pos = int(key)
|
||||
result[pos] = _apply(result[pos], value)
|
||||
return result
|
||||
|
||||
result = dict(state) if isinstance(state, dict) else {}
|
||||
for key, value in diff.items():
|
||||
if key == "$replace":
|
||||
return value
|
||||
if key == "$delete":
|
||||
if isinstance(value, list):
|
||||
for k in value:
|
||||
result.pop(k, None)
|
||||
else:
|
||||
result.pop(value, None)
|
||||
keys = value if isinstance(value, list) else [value]
|
||||
for k in keys:
|
||||
result.pop(k, None)
|
||||
elif key == "$insert":
|
||||
continue
|
||||
if isinstance(value, dict):
|
||||
old = result.get(key, {})
|
||||
if not isinstance(old, dict):
|
||||
old = {}
|
||||
result[key] = _apply_diff(old, value)
|
||||
continue
|
||||
result[key] = value
|
||||
|
||||
elif key in result:
|
||||
result[key] = _apply(result[key], value)
|
||||
else:
|
||||
result[key] = value
|
||||
return result
|
||||
|
||||
@@ -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
|
||||
|
||||
|
||||
@@ -88,6 +88,11 @@ def transaction(
|
||||
new_dict = struct_to_dict(impl.data, serializer=impl.serializer)
|
||||
diff = compute_diff(impl.statedict, new_dict)
|
||||
if diff:
|
||||
if impl.callback_registry.has("validate"):
|
||||
impl.callback_registry.invoke_sync(
|
||||
"validate",
|
||||
InjectionContext(data=impl.data, kanta=impl._kanta),
|
||||
)
|
||||
previous = impl.statedict
|
||||
record = impl.queue_change(action, new_dict, user=user, mtime=mtime)
|
||||
if record is not None:
|
||||
|
||||
+2
-1
@@ -17,7 +17,6 @@ readme = "README.md"
|
||||
requires-python = ">=3.11"
|
||||
dependencies = [
|
||||
"blake3>=1.0.8",
|
||||
"jsondiff>=2.2.1",
|
||||
"msgspec>=0.20.0",
|
||||
]
|
||||
|
||||
@@ -30,10 +29,12 @@ bin = [
|
||||
]
|
||||
|
||||
[project.urls]
|
||||
Homepage = "https://vasanko.com/coders/kanta"
|
||||
Repository = "https://git.zi.fi/LeoVasanko/kanta"
|
||||
|
||||
[dependency-groups]
|
||||
dev = [
|
||||
"jsondiff>=2.2.1",
|
||||
"pytest>=9.0.2",
|
||||
"pytest-asyncio>=1.3.0",
|
||||
]
|
||||
|
||||
+243
-1
@@ -1,8 +1,26 @@
|
||||
from kanta.diff import compute_diff
|
||||
"""Tests for our own diff producer/consumer and jsondiff compatibility.
|
||||
|
||||
jsondiff is a dev dependency used only here, to verify that:
|
||||
|
||||
- jsondiff.patch(..., marshal=True) can apply patches produced by
|
||||
compute_diff (our format is a subset of jsondiff's marshaled syntax);
|
||||
- apply_diff can apply patches produced by jsondiff.diff(..., marshal=True),
|
||||
including positional $insert/$delete list edits and per-index nested diffs.
|
||||
"""
|
||||
|
||||
import jsondiff
|
||||
import pytest
|
||||
|
||||
from kanta.diff import compute_diff, patch_state
|
||||
from kanta.logging import format_diff
|
||||
from kanta.serialization.base import apply_diff
|
||||
|
||||
# --- Producer: compute_diff ------------------------------------------------
|
||||
|
||||
|
||||
def test_no_diff():
|
||||
assert compute_diff({"a": 1}, {"a": 1}) is None
|
||||
assert compute_diff({}, {}) is None
|
||||
|
||||
|
||||
def test_simple_diff():
|
||||
@@ -14,3 +32,227 @@ def test_simple_diff():
|
||||
def test_nested_diff():
|
||||
diff = compute_diff({"x": {"y": 1}}, {"x": {"y": 2}})
|
||||
assert diff == {"x": {"y": 2}}
|
||||
|
||||
|
||||
def test_key_added():
|
||||
assert compute_diff({"a": 1}, {"a": 1, "b": 2}) == {"b": 2}
|
||||
|
||||
|
||||
def test_key_removed():
|
||||
assert compute_diff({"a": 1, "b": 2}, {"a": 1}) == {"$delete": ["b"]}
|
||||
|
||||
|
||||
def test_last_key_removed_is_delete_not_replace():
|
||||
# jsondiff's minimal-diff search emits {"$replace": {}} here; we emit
|
||||
# what actually happened: the key was deleted.
|
||||
assert compute_diff({"a": 1}, {}) == {"$delete": ["a"]}
|
||||
assert compute_diff({"x": {"y": 1}}, {"x": {}}) == {"x": {"$delete": ["y"]}}
|
||||
|
||||
|
||||
def test_list_changes_are_full_assignment():
|
||||
# No $insert/$delete positional edits: lists are replaced wholesale.
|
||||
assert compute_diff({"l": [1, 2]}, {"l": [1, 2, 3]}) == {"l": [1, 2, 3]}
|
||||
assert compute_diff({"l": [1, 2, 3]}, {"l": [1, 3]}) == {"l": [1, 3]}
|
||||
assert compute_diff({"l": [1]}, {"l": []}) == {"l": []}
|
||||
|
||||
|
||||
def test_list_with_unchanged_prefix_is_full_assignment():
|
||||
diff = compute_diff({"l": ["a", "b", "c"]}, {"l": ["a", "x", "b", "c"]})
|
||||
assert diff == {"l": ["a", "x", "b", "c"]}
|
||||
|
||||
|
||||
def test_type_changes_are_full_assignment():
|
||||
assert compute_diff({"a": {"x": 1}}, {"a": [1]}) == {"a": [1]}
|
||||
# A dict replacing a non-dict is a plain assignment too: the consumer
|
||||
# sees from the old value whether to patch (dict) or replace.
|
||||
assert compute_diff({"a": [1]}, {"a": {"x": 1}}) == {"a": {"x": 1}}
|
||||
assert compute_diff({"a": 1}, {"a": None}) == {"a": None}
|
||||
|
||||
|
||||
def test_new_dict_value_assigned_wholesale():
|
||||
assert compute_diff({}, {"a": {"x": 1}}) == {"a": {"x": 1}}
|
||||
|
||||
|
||||
def test_dollar_keys_escaped():
|
||||
assert compute_diff({}, {"$weird": 1}) == {"$$weird": 1}
|
||||
assert compute_diff({"$weird": 1}, {"$weird": 2}) == {"$$weird": 2}
|
||||
assert compute_diff({"$weird": 1}, {}) == {"$delete": ["$$weird"]}
|
||||
|
||||
|
||||
def test_dollar_values_not_escaped():
|
||||
# Only keys are escaped; values are stored verbatim, even "$delete".
|
||||
assert compute_diff({"s": 1}, {"s": "$y"}) == {"s": "$y"}
|
||||
assert compute_diff({"s": 1}, {"s": "$delete"}) == {"s": "$delete"}
|
||||
assert compute_diff({}, {"o": {"s": "$y", "l": ["$z"]}}) == {
|
||||
"o": {"s": "$y", "l": ["$z"]}
|
||||
}
|
||||
|
||||
|
||||
# --- Consumer: apply_diff / patch_state -------------------------------------
|
||||
|
||||
|
||||
def test_patch_state_delegates():
|
||||
assert patch_state({"a": 1}, {"a": 2}) == {"a": 2}
|
||||
|
||||
|
||||
def test_apply_scalar_and_add():
|
||||
assert apply_diff({"a": 1}, {"a": 2, "b": 3}) == {"a": 2, "b": 3}
|
||||
|
||||
|
||||
def test_apply_delete():
|
||||
assert apply_diff({"a": 1, "b": 2}, {"$delete": ["b"]}) == {"a": 1}
|
||||
assert apply_diff({"a": 1}, {"$delete": ["a"]}) == {}
|
||||
|
||||
|
||||
def test_apply_replace():
|
||||
assert apply_diff({"a": 1, "b": 2}, {"$replace": {"c": 3}}) == {"c": 3}
|
||||
assert apply_diff({"x": {"a": 1}}, {"x": {"$replace": [1]}}) == {"x": [1]}
|
||||
|
||||
|
||||
def test_apply_nested_delete():
|
||||
diff = {"x": {"$delete": ["y"]}}
|
||||
assert apply_diff({"x": {"y": 2, "z": 3}}, diff) == {"x": {"z": 3}}
|
||||
|
||||
|
||||
def test_apply_list_insert():
|
||||
diff = {"l": {"$insert": [[1, "x"]]}}
|
||||
assert apply_diff({"l": ["a", "b"]}, diff) == {"l": ["a", "x", "b"]}
|
||||
|
||||
|
||||
def test_apply_list_delete():
|
||||
diff = {"l": {"$delete": [1]}}
|
||||
assert apply_diff({"l": ["a", "b", "c"]}, diff) == {"l": ["a", "c"]}
|
||||
|
||||
|
||||
def test_apply_list_delete_multiple_positions():
|
||||
# jsondiff emits positions in descending order for sequential pops.
|
||||
diff = {"l": {"$delete": [4, 2, 0]}}
|
||||
assert apply_diff({"l": [0, 1, 2, 3, 4]}, diff) == {"l": [1, 3]}
|
||||
|
||||
|
||||
def test_apply_list_insert_and_delete():
|
||||
diff = {"l": {"$insert": [[0, 9], [2, 8], [4, 7]], "$delete": [2, 0]}}
|
||||
assert apply_diff({"l": [0, 1, 2, 3]}, diff) == {"l": [9, 1, 8, 3, 7]}
|
||||
|
||||
|
||||
def test_apply_list_per_index_nested_diff():
|
||||
diff = {"l": {"1": {"y": 3}}}
|
||||
state = {"l": [{"x": 1}, {"y": 2}]}
|
||||
assert apply_diff(state, diff) == {"l": [{"x": 1}, {"y": 3}]}
|
||||
|
||||
|
||||
def test_apply_escaped_keys_and_values():
|
||||
assert apply_diff({}, {"$$weird": 1}) == {"$weird": 1}
|
||||
assert apply_diff({"$weird": 1}, {"$delete": ["$$weird"]}) == {}
|
||||
# jsondiff escapes $-values as "$$.."; those are unescaped on apply.
|
||||
assert apply_diff({"s": 1}, {"s": "$$y"}) == {"s": "$y"}
|
||||
assert apply_diff({"s": 1}, {"s": "$$delete"}) == {"s": "$delete"}
|
||||
# Our own producer stores values verbatim; single-$ stays as-is.
|
||||
assert apply_diff({"s": 1}, {"s": "$y"}) == {"s": "$y"}
|
||||
assert apply_diff({"s": 1}, {"s": "$delete"}) == {"s": "$delete"}
|
||||
assert apply_diff({}, {"o": {"s": "$$y", "l": ["$$z"]}}) == {
|
||||
"o": {"s": "$y", "l": ["$z"]}
|
||||
}
|
||||
|
||||
|
||||
def test_apply_empty_diff():
|
||||
assert apply_diff({"a": 1}, {}) == {"a": 1}
|
||||
|
||||
|
||||
def test_apply_diff_on_missing_state():
|
||||
assert apply_diff({}, {"a": {"b": 1}}) == {"a": {"b": 1}}
|
||||
|
||||
|
||||
def test_apply_bare_dict_replaces_non_dict():
|
||||
# Our own producer emits no $replace; a dict over a non-dict old value
|
||||
# is a wholesale replacement.
|
||||
assert apply_diff({"a": [1, 2]}, {"a": {"x": 1}}) == {"a": {"x": 1}}
|
||||
assert apply_diff({"a": 5}, {"a": {"x": 1}}) == {"a": {"x": 1}}
|
||||
assert apply_diff({"a": None}, {"a": {"x": 1}}) == {"a": {"x": 1}}
|
||||
|
||||
|
||||
def test_apply_list_patch_still_works_on_lists():
|
||||
# jsondiff-style per-index diff keeps list-op semantics on list state.
|
||||
assert apply_diff({"l": [1, 2]}, {"l": {"1": 9}}) == {"l": [1, 9]}
|
||||
|
||||
|
||||
# --- jsondiff compatibility, both directions --------------------------------
|
||||
|
||||
COMPAT_CASES = [
|
||||
("scalar change", {"a": 1}, {"a": 2}),
|
||||
("key add", {"a": 1}, {"a": 1, "b": 2}),
|
||||
("key remove", {"a": 1, "b": 2}, {"a": 1}),
|
||||
("last key removed", {"a": 1}, {}),
|
||||
("nested delete", {"a": {"x": 1, "y": 2}}, {"a": {"x": 1}}),
|
||||
("nested mixed", {"a": {"x": 1, "y": 2}}, {"a": {"x": 9, "z": 3}}),
|
||||
("list append", {"l": [1, 2]}, {"l": [1, 2, 3]}),
|
||||
("list insert mid", {"l": [1, 2, 3]}, {"l": [1, 9, 2, 3]}),
|
||||
("list remove mid", {"l": [1, 2, 3]}, {"l": [1, 3]}),
|
||||
("list remove many", {"l": [0, 1, 2, 3, 4]}, {"l": [1, 3]}),
|
||||
("list replace all", {"l": [1, 2]}, {"l": [3, 4]}),
|
||||
("list insert+delete", {"l": [0, 1, 2, 3]}, {"l": [9, 1, 8, 3, 7]}),
|
||||
("dict in list", {"l": [{"x": 1}, {"y": 2}]}, {"l": [{"x": 1}, {"y": 3}]}),
|
||||
("list to empty", {"l": [1]}, {"l": []}),
|
||||
("type change dict->list", {"a": {"x": 1}}, {"a": [1]}),
|
||||
("type change list->dict", {"a": [1]}, {"a": {"x": 1}}),
|
||||
("dollar key", {"$k": 1, "b": 1}, {"$k": 2}),
|
||||
("dollar value", {"s": "$x"}, {"s": "$y"}),
|
||||
(
|
||||
"deep nesting",
|
||||
{"a": {"b": {"c": {"d": 1, "e": 2}}}},
|
||||
{"a": {"b": {"c": {"d": 9}}}},
|
||||
),
|
||||
]
|
||||
|
||||
|
||||
# jsondiff.patch cannot apply our patches for "type change list->dict"
|
||||
# (we emit a bare dict where jsondiff needs $replace) and "dollar value"
|
||||
# (we do not escape "$"-prefixed values; jsondiff.patch would strip the
|
||||
# "$"), so those cases are excluded from this direction.
|
||||
JSONDIFF_APPLIES_CASES = [
|
||||
c for c in COMPAT_CASES if c[0] not in {"type change list->dict", "dollar value"}
|
||||
]
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
"name,old,new", JSONDIFF_APPLIES_CASES, ids=[c[0] for c in JSONDIFF_APPLIES_CASES]
|
||||
)
|
||||
def test_jsondiff_applies_our_patches(name, old, new):
|
||||
diff = compute_diff(old, new)
|
||||
assert diff is not None
|
||||
assert jsondiff.patch(old, diff, marshal=True) == new
|
||||
|
||||
|
||||
@pytest.mark.parametrize("name,old,new", COMPAT_CASES, ids=[c[0] for c in COMPAT_CASES])
|
||||
def test_we_apply_jsondiff_patches(name, old, new):
|
||||
diff = jsondiff.diff(old, new, marshal=True)
|
||||
assert apply_diff(old, diff) == new
|
||||
|
||||
|
||||
@pytest.mark.parametrize("name,old,new", COMPAT_CASES, ids=[c[0] for c in COMPAT_CASES])
|
||||
def test_our_own_round_trip(name, old, new):
|
||||
diff = compute_diff(old, new)
|
||||
assert diff is not None
|
||||
assert apply_diff(old, diff) == new
|
||||
|
||||
|
||||
def test_no_diff_means_equal_states():
|
||||
for _name, old, new in COMPAT_CASES:
|
||||
assert compute_diff(old, new) is not None # cases really differ
|
||||
assert compute_diff({"a": [1, {"b": "$x"}]}, {"a": [1, {"b": "$x"}]}) is None
|
||||
|
||||
|
||||
# --- Logging ----------------------------------------------------------------
|
||||
|
||||
|
||||
def test_format_diff_list_edit_shows_whole_list():
|
||||
diff = jsondiff.diff({"l": [1, 2, 3]}, {"l": [1, 9, 3]}, marshal=True)
|
||||
lines = format_diff(diff, previous={"l": [1, 2, 3]})
|
||||
text = "\n".join(lines)
|
||||
assert "$insert" not in text
|
||||
assert "[1, 9, 3]" in text
|
||||
|
||||
|
||||
def test_format_diff_unescapes_dollar_keys():
|
||||
lines = format_diff({"$$weird": 1}, previous={})
|
||||
assert any("$weird" in line and "$$weird" not in line for line in lines)
|
||||
|
||||
@@ -0,0 +1,55 @@
|
||||
"""Tests for LockedFile low-level behaviors."""
|
||||
|
||||
from kanta.exceptions import FileLockError
|
||||
from kanta.filelock import LockedFile
|
||||
|
||||
|
||||
def test_replace_content_rewrites_in_place(tmp_path):
|
||||
path = tmp_path / "data.kantadb"
|
||||
path.write_bytes(b"original content here")
|
||||
|
||||
f = LockedFile()
|
||||
f.open(path)
|
||||
try:
|
||||
f.replace_content(b"new")
|
||||
assert f.size() == 3
|
||||
f.write(b"!")
|
||||
finally:
|
||||
f.close()
|
||||
|
||||
assert path.read_bytes() == b"new!"
|
||||
|
||||
|
||||
def test_replace_content_keeps_lock(tmp_path):
|
||||
path = tmp_path / "data.kantadb"
|
||||
path.write_bytes(b"abc")
|
||||
|
||||
f = LockedFile()
|
||||
f.open(path)
|
||||
try:
|
||||
f.replace_content(b"xyz")
|
||||
other = LockedFile()
|
||||
try:
|
||||
other.open(path)
|
||||
raise AssertionError("second open should fail while lock is held")
|
||||
except FileLockError:
|
||||
pass
|
||||
finally:
|
||||
f.close()
|
||||
|
||||
|
||||
def test_replace_content_grow_and_shrink(tmp_path):
|
||||
path = tmp_path / "data.kantadb"
|
||||
path.write_bytes(b"x" * 100)
|
||||
|
||||
f = LockedFile()
|
||||
f.open(path)
|
||||
try:
|
||||
f.replace_content(b"")
|
||||
assert f.size() == 0
|
||||
f.replace_content(b"y" * 200)
|
||||
assert f.size() == 200
|
||||
finally:
|
||||
f.close()
|
||||
|
||||
assert path.read_bytes() == b"y" * 200
|
||||
@@ -47,7 +47,7 @@ async def test_new_file_writes_bootstrap_record_without_handlers(
|
||||
records = read_changes(path, format_config)
|
||||
assert len(records) == 1
|
||||
assert records[0].a == "bootstrap"
|
||||
assert records[0].diff == {"$replace": {"users": {}, "counter": 0}}
|
||||
assert records[0].diff == {"users": {}, "counter": 0}
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
@@ -63,7 +63,8 @@ async def test_new_file_persists_initial_state_for_roundtrip(tmp_path, format_co
|
||||
assert len(records) == 1
|
||||
assert records[0].a == "bootstrap"
|
||||
assert records[0].diff == {
|
||||
"$replace": {"users": {"alice": {"name": "Alice", "age": 0}}, "counter": 5}
|
||||
"users": {"alice": {"name": "Alice", "age": 0}},
|
||||
"counter": 5,
|
||||
}
|
||||
|
||||
kanta2 = make_kanta(path, Data, format_config)
|
||||
|
||||
@@ -183,7 +183,7 @@ def test_apply_returns_change_information():
|
||||
assert result.migrations[0].name == "migrate_v1"
|
||||
assert result.migrations[0].description == "Set x"
|
||||
assert result.migrations[0].changed is True
|
||||
assert result.migrations[0].diff == {"$replace": {"x": 1}}
|
||||
assert result.migrations[0].diff == {"x": 1}
|
||||
|
||||
assert result.migrations[1].name == "migrate_v2"
|
||||
assert result.migrations[1].description == "No-op"
|
||||
|
||||
@@ -0,0 +1,214 @@
|
||||
"""Tests for retention-based database rotation (docs/rotation.md)."""
|
||||
|
||||
from datetime import UTC, datetime, timedelta
|
||||
from pathlib import Path
|
||||
|
||||
import pytest
|
||||
|
||||
from kanta.structs import ChangeRecord, Snapshot
|
||||
from tests.support import Data, make_kanta, read_changes
|
||||
|
||||
pytestmark = pytest.mark.asyncio
|
||||
|
||||
DAY = timedelta(days=1)
|
||||
T0 = datetime(2026, 1, 1, tzinfo=UTC)
|
||||
|
||||
|
||||
def make_clock(cell: list[datetime]):
|
||||
def clock() -> datetime:
|
||||
return cell[0]
|
||||
|
||||
return clock
|
||||
|
||||
|
||||
async def write_history(path: Path, format_config, days: list[int]) -> None:
|
||||
"""Write one change per day offset (relative to T0) with a fake clock."""
|
||||
cell = [T0 + (days[0] - 1) * DAY] # bootstrap predates all history
|
||||
kanta = make_kanta(path, Data, format_config)
|
||||
kanta.clock(make_clock(cell))
|
||||
await kanta.open(log=False)
|
||||
for day in days:
|
||||
cell[0] = T0 + day * DAY
|
||||
with kanta.transaction(f"day{day}", log=False) as data:
|
||||
data.counter += 1
|
||||
await kanta.flush()
|
||||
await kanta.close()
|
||||
|
||||
|
||||
def read_all(path: Path, format_config):
|
||||
"""All records (changes and snapshots) in file order."""
|
||||
_, serializer_cls = format_config
|
||||
serializer = serializer_cls()
|
||||
framer = serializer.framer_cls()
|
||||
out = []
|
||||
for is_snapshot, payload, _, _ in framer.iter_records(path.read_bytes(), 0):
|
||||
out.append(
|
||||
serializer.decode(payload, type=Snapshot if is_snapshot else ChangeRecord)
|
||||
)
|
||||
return out
|
||||
|
||||
|
||||
def rotated_files(path: Path) -> list[Path]:
|
||||
return sorted(path.parent.glob(f"{path.stem}@*.kantadb"))
|
||||
|
||||
|
||||
async def test_rotation_splits_history(tmp_path, format_config):
|
||||
path = tmp_path / "data.kantadb"
|
||||
await write_history(path, format_config, days=[-40, -20, -5])
|
||||
|
||||
cell = [T0]
|
||||
kanta = make_kanta(path, Data, format_config, retention=30 * DAY)
|
||||
kanta.clock(make_clock(cell))
|
||||
await kanta.open(log=False)
|
||||
assert kanta.data.counter == 3
|
||||
await kanta.close()
|
||||
|
||||
rotated = rotated_files(path)
|
||||
assert len(rotated) == 1
|
||||
|
||||
# Main file: leading snapshot (ts = last dropped record), the retained
|
||||
# changes, and no final snapshot (too few retained changes).
|
||||
records = read_all(path, format_config)
|
||||
assert isinstance(records[0], Snapshot)
|
||||
assert records[0].ts == T0 - 40 * DAY
|
||||
assert records[0].state["counter"] == 1
|
||||
changes = [r for r in records if isinstance(r, ChangeRecord)]
|
||||
assert [c.a for c in changes] == ["day-20", "day-5"]
|
||||
|
||||
# Rotated file holds exactly the dropped history, ending at the last
|
||||
# dropped record whose ts matches the filename.
|
||||
stamp = (T0 - 40 * DAY).strftime("%Y%m%dT%H%M%SZ")
|
||||
assert rotated[0].name == f"data@{stamp}.kantadb"
|
||||
rrecords = read_all(rotated[0], format_config)
|
||||
assert [r.a for r in rrecords] == ["bootstrap", "day-40"]
|
||||
|
||||
|
||||
async def test_rotation_reopens_cleanly_and_does_not_rerotate(tmp_path, format_config):
|
||||
path = tmp_path / "data.kantadb"
|
||||
await write_history(path, format_config, days=[-40, -5])
|
||||
|
||||
cell = [T0]
|
||||
for expected_changes in (["day-5"], ["day-5"]):
|
||||
kanta = make_kanta(path, Data, format_config, retention=30 * DAY)
|
||||
kanta.clock(make_clock(cell))
|
||||
async with kanta:
|
||||
assert kanta.data.counter == 2
|
||||
assert [c.a for c in read_changes(path, format_config)] == expected_changes
|
||||
|
||||
# Second open found a file whose history already fits the window.
|
||||
assert len(rotated_files(path)) == 1
|
||||
|
||||
|
||||
async def test_rotation_noop_when_retention_covers_all(tmp_path, format_config):
|
||||
path = tmp_path / "data.kantadb"
|
||||
await write_history(path, format_config, days=[-5])
|
||||
before = path.read_bytes()
|
||||
|
||||
cell = [T0]
|
||||
kanta = make_kanta(path, Data, format_config, retention=30 * DAY)
|
||||
kanta.clock(make_clock(cell))
|
||||
async with kanta:
|
||||
assert kanta.data.counter == 1
|
||||
|
||||
assert rotated_files(path) == []
|
||||
assert path.read_bytes() == before
|
||||
|
||||
|
||||
async def test_rotation_aged_out_database_reduces_to_single_snapshot(
|
||||
tmp_path, format_config
|
||||
):
|
||||
path = tmp_path / "data.kantadb"
|
||||
await write_history(path, format_config, days=[-40, -35])
|
||||
|
||||
cell = [T0]
|
||||
kanta = make_kanta(path, Data, format_config, retention=30 * DAY)
|
||||
kanta.clock(make_clock(cell))
|
||||
async with kanta:
|
||||
assert kanta.data.counter == 2
|
||||
|
||||
records = read_all(path, format_config)
|
||||
assert len(records) == 1
|
||||
assert isinstance(records[0], Snapshot)
|
||||
assert records[0].state["counter"] == 2
|
||||
|
||||
# Opening again must not rotate the snapshot-only file.
|
||||
before = path.read_bytes()
|
||||
kanta = make_kanta(path, Data, format_config, retention=30 * DAY)
|
||||
kanta.clock(make_clock(cell))
|
||||
async with kanta:
|
||||
assert kanta.data.counter == 2
|
||||
assert path.read_bytes() == before
|
||||
assert len(rotated_files(path)) == 1
|
||||
|
||||
|
||||
async def test_rotation_validates_against_internal_snapshots(tmp_path, format_config):
|
||||
path = tmp_path / "data.kantadb"
|
||||
cell = [T0 - 40 * DAY]
|
||||
kanta = make_kanta(path, Data, format_config)
|
||||
kanta.clock(make_clock(cell))
|
||||
await kanta.open(log=False)
|
||||
with kanta.transaction("old", log=False) as data:
|
||||
data.counter = 1
|
||||
await kanta.flush()
|
||||
kanta.request_snapshot()
|
||||
kanta._impl.maybe_snapshot()
|
||||
cell[0] = T0 - 1 * DAY
|
||||
with kanta.transaction("new", log=False) as data:
|
||||
data.counter = 2
|
||||
await kanta.flush()
|
||||
await kanta.close()
|
||||
|
||||
cell[0] = T0
|
||||
kanta = make_kanta(path, Data, format_config, retention=30 * DAY)
|
||||
kanta.clock(make_clock(cell))
|
||||
async with kanta:
|
||||
assert kanta.data.counter == 2
|
||||
|
||||
records = read_all(path, format_config)
|
||||
assert isinstance(records[0], Snapshot)
|
||||
assert records[0].state["counter"] == 1
|
||||
assert [r.a for r in records if isinstance(r, ChangeRecord)] == ["new"]
|
||||
|
||||
|
||||
@pytest.mark.parametrize("name", ["data", "data.db", "data.kantadb"])
|
||||
async def test_rotated_naming_normalizes_extension(tmp_path, format_config, name):
|
||||
path = tmp_path / name
|
||||
await write_history(path, format_config, days=[-40, -5])
|
||||
|
||||
cell = [T0]
|
||||
kanta = make_kanta(path, Data, format_config, retention=30 * DAY)
|
||||
kanta.clock(make_clock(cell))
|
||||
async with kanta:
|
||||
pass
|
||||
|
||||
stamp = (T0 - 40 * DAY).strftime("%Y%m%dT%H%M%SZ")
|
||||
assert (tmp_path / f"data@{stamp}.kantadb").exists()
|
||||
|
||||
|
||||
async def test_retention_accepts_int_days(tmp_path, format_config):
|
||||
path = tmp_path / "data.kantadb"
|
||||
await write_history(path, format_config, days=[-40, -5])
|
||||
|
||||
cell = [T0]
|
||||
kanta = make_kanta(path, Data, format_config, retention=30)
|
||||
kanta.clock(make_clock(cell))
|
||||
async with kanta:
|
||||
assert kanta.data.counter == 2
|
||||
|
||||
assert len(rotated_files(path)) == 1
|
||||
assert [c.a for c in read_changes(path, format_config)] == ["day-5"]
|
||||
|
||||
|
||||
async def test_rotation_disabled_by_default(tmp_path, format_config):
|
||||
path = tmp_path / "data.kantadb"
|
||||
await write_history(path, format_config, days=[-40, -5])
|
||||
before = path.read_bytes()
|
||||
|
||||
cell = [T0]
|
||||
kanta = make_kanta(path, Data, format_config)
|
||||
kanta.clock(make_clock(cell))
|
||||
async with kanta:
|
||||
assert kanta.data.counter == 2
|
||||
|
||||
assert rotated_files(path) == []
|
||||
assert path.read_bytes() == before
|
||||
@@ -0,0 +1,116 @@
|
||||
"""Tests for the @kanta.validate integrity-validation callbacks."""
|
||||
|
||||
import pytest
|
||||
|
||||
from tests.support import Data, make_kanta
|
||||
|
||||
pytestmark = pytest.mark.asyncio
|
||||
|
||||
|
||||
async def test_validate_passes_on_valid_data(tmp_path, format_config):
|
||||
kanta = make_kanta(tmp_path / "d.kantadb", Data, format_config)
|
||||
calls = []
|
||||
|
||||
@kanta.validate
|
||||
def check(data: Data):
|
||||
calls.append(data.counter)
|
||||
assert data.counter >= 0
|
||||
|
||||
async with kanta:
|
||||
with kanta.transaction("inc", log=False) as data:
|
||||
data.counter = 1
|
||||
|
||||
assert calls # ran during bootstrap/open and the transaction
|
||||
|
||||
|
||||
async def test_validate_failure_rolls_back_transaction(tmp_path, format_config):
|
||||
kanta = make_kanta(tmp_path / "d.kantadb", Data, format_config)
|
||||
|
||||
@kanta.validate
|
||||
def check(data: Data):
|
||||
if data.counter < 0:
|
||||
raise ValueError("counter must not go negative")
|
||||
|
||||
await kanta.open(log=False)
|
||||
with pytest.raises(ValueError, match="negative"):
|
||||
with kanta.transaction("dec", log=False) as data:
|
||||
data.counter = -1
|
||||
assert kanta.data.counter == 0 # rolled back
|
||||
await kanta.close()
|
||||
|
||||
# The invalid change never reached the history.
|
||||
kanta2 = make_kanta(tmp_path / "d.kantadb", Data, format_config)
|
||||
async with kanta2:
|
||||
assert kanta2.data.counter == 0
|
||||
|
||||
|
||||
async def test_validate_runs_on_open_after_replay(tmp_path, format_config):
|
||||
kanta = make_kanta(tmp_path / "d.kantadb", Data, format_config)
|
||||
async with kanta:
|
||||
with kanta.transaction("set", log=False) as data:
|
||||
data.counter = 5
|
||||
|
||||
kanta2 = make_kanta(tmp_path / "d.kantadb", Data, format_config)
|
||||
seen = []
|
||||
|
||||
@kanta2.validate
|
||||
def check(data: Data):
|
||||
seen.append(data.counter)
|
||||
|
||||
async with kanta2:
|
||||
pass
|
||||
assert 5 in seen
|
||||
|
||||
|
||||
async def test_validate_failure_aborts_open(tmp_path, format_config):
|
||||
kanta = make_kanta(tmp_path / "d.kantadb", Data, format_config)
|
||||
async with kanta:
|
||||
with kanta.transaction("set", log=False) as data:
|
||||
data.counter = 5
|
||||
|
||||
kanta2 = make_kanta(tmp_path / "d.kantadb", Data, format_config)
|
||||
|
||||
@kanta2.validate
|
||||
def check(data: Data):
|
||||
raise ValueError("always inconsistent")
|
||||
|
||||
with pytest.raises(ValueError, match="inconsistent"):
|
||||
await kanta2.open(log=False)
|
||||
|
||||
# The failed open released the file: a fresh instance can open it.
|
||||
kanta3 = make_kanta(tmp_path / "d.kantadb", Data, format_config)
|
||||
async with kanta3:
|
||||
assert kanta3.data.counter == 5
|
||||
|
||||
|
||||
async def test_multiple_validators_stop_at_first_failure(tmp_path, format_config):
|
||||
kanta = make_kanta(tmp_path / "d.kantadb", Data, format_config)
|
||||
calls = []
|
||||
|
||||
@kanta.validate
|
||||
def first(data: Data):
|
||||
calls.append("first")
|
||||
if data.counter > 1:
|
||||
raise ValueError("too big")
|
||||
|
||||
@kanta.validate
|
||||
def second(data: Data):
|
||||
calls.append("second")
|
||||
|
||||
await kanta.open(log=False)
|
||||
calls.clear()
|
||||
with pytest.raises(ValueError, match="too big"):
|
||||
with kanta.transaction("bump", log=False) as data:
|
||||
data.counter = 2
|
||||
assert calls == ["first"]
|
||||
await kanta.close()
|
||||
|
||||
|
||||
async def test_validate_rejects_async_callback(tmp_path, format_config):
|
||||
kanta = make_kanta(tmp_path / "d.kantadb", Data, format_config)
|
||||
|
||||
with pytest.raises(TypeError, match="must not be async"):
|
||||
|
||||
@kanta.validate
|
||||
async def check(data: Data):
|
||||
pass
|
||||
Reference in New Issue
Block a user