diff --git a/README.md b/README.md index 0176dcf..89270bf 100644 --- a/README.md +++ b/README.md @@ -12,18 +12,18 @@ you `delete` or `clear` them. `requirements.txt`: ``` -aiokv @ git+ssh://git@git.rethinkstudios.io/rethink-public/aiokv.git@v0.1.1 +aiokv @ git+ssh://git@git.rethinkstudios.io/rethink-public/aiokv.git@v0.2.0 ``` Direct: ```bash -pip install "aiokv @ git+ssh://git@git.rethinkstudios.io/rethink-public/aiokv.git@v0.1.1" +pip install "aiokv @ git+ssh://git@git.rethinkstudios.io/rethink-public/aiokv.git@v0.2.0" ``` Requires `aiofiles` (pulled transitively). -Drop the `@v0.1.1` suffix from the line above to install the latest unpinned. +Drop the `@v0.2.0` suffix from the line above to install the latest unpinned. ## Usage @@ -62,14 +62,15 @@ Prefer `AioKV` in new code. ## Durability Writes are **atomic**: data is written to a temp file in the same directory and -`os.replace()`d over the target (atomic on POSIX). A **process** crash mid-write leaves -the previous good file intact, and a reader never observes a partial file. (This is -process-crash safety, not power-loss durability — there's no `fsync`, so an OS/power -failure could still lose the last write; fine for reconstructible single-process state.) -A single -`asyncio.Lock` guards every read and write, so concurrent operations on one instance -are consistent and no update is lost. All blocking filesystem calls run via -`asyncio.to_thread`, so nothing stalls the event loop. +`os.replace()`d over the target's realpath — symlink-safe, so a symlinked store file is +written through rather than clobbered. A **process** crash mid-write leaves the previous +good file intact, and a reader never observes a partial file. (This is process-crash +safety, not power-loss durability — there's no `fsync`, so an OS/power failure could +still lose the last write; fine for reconstructible single-process state.) Orphaned +`..tmp` files left by a hard crash of a *different, dead* process are swept on the +next save. A single `asyncio.Lock` guards every read and write, so concurrent operations +on one instance are consistent and no update is lost. All blocking filesystem calls and +JSON (de)serialization run via `asyncio.to_thread`, so nothing stalls the event loop. ## Scope — read this @@ -82,10 +83,13 @@ are consistent and no update is lost. All blocking filesystem calls run via ## Error contract +- Keys must be `str` — `get` / `set` raise `ValueError` on a non-str key (JSON object + keys are always strings, so a non-str key would silently never round-trip). - `get` / `set` / `get_all` raise on unexpected I/O. `_load` raises `JSONDecodeError` on a truncated/corrupt file, and `ValueError` when the file holds valid JSON that isn't an object (a bare list/number/string/null) — so corruption or a wrong-shaped - file is visible rather than silently masked. + file is visible rather than silently masked. Non-finite floats (`NaN`/`Infinity`) + raise `ValueError` on `set` rather than persisting invalid JSON. - `delete` / `clear` log the exception and return `False` on error, `True` otherwise. ## Versioning diff --git a/pyproject.toml b/pyproject.toml index 10d259f..d9a5d51 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -4,7 +4,7 @@ build-backend = "hatchling.build" [project] name = "aiokv" -version = "0.1.1" +version = "0.2.0" description = "Async file-backed key-value store for single-process local state — atomic writes, no TTL, config-free, installable." requires-python = ">=3.10" dependencies = [ diff --git a/src/aiokv/__init__.py b/src/aiokv/__init__.py index 8c9624f..67b5f7d 100644 --- a/src/aiokv/__init__.py +++ b/src/aiokv/__init__.py @@ -1,3 +1,5 @@ from .aiokv import AioKV, aiocache +__version__ = "0.2.0" + __all__ = ["AioKV", "aiocache"] diff --git a/src/aiokv/aiokv.py b/src/aiokv/aiokv.py index 061d3ea..71aeeca 100644 --- a/src/aiokv/aiokv.py +++ b/src/aiokv/aiokv.py @@ -13,10 +13,12 @@ eviction — values live until you delete or clear them. when = await kv.get("ran_cleanup") await kv.delete("last_seen") -durability: writes are atomic — data is written to a temp file in the same -directory and os.replace()d over the target, so a crash mid-write never corrupts -the store and readers never see a partial file. a single asyncio.Lock guards every -read and write, so concurrent operations on one instance are consistent. +durability: atomic writes (temp file + os.replace, symlink-safe), stale +`..tmp` files from a hard crash are swept on save. this is process-crash +safety, NOT power-loss durability — no fsync, so an OS/power failure can still +lose the last write. a single asyncio.Lock guards every read and write, so +concurrent operations on one instance are consistent. JSON (de)serialization +runs off-loop via asyncio.to_thread so a large store never stalls other tasks. scope: SINGLE-PROCESS, single-instance local state only. the lock is per-instance — two AioKV instances (or two processes) pointing at the same file are NOT safe and @@ -24,13 +26,16 @@ will clobber each other. for shared cross-process/cross-bot state, use a databas (e.g. mongo), not this. config-free: the file path is passed at construction; nothing is read from a global -config. errors in delete/clear are logged and swallowed (returning False); get/set -raise on unexpected i/o so a real failure is visible. +config. keys must be str (ValueError otherwise — JSON object keys are always +strings, so a non-str key would silently never round-trip). errors in delete/clear +are logged and swallowed (returning False); get/set raise on unexpected i/o so a +real failure is visible. """ import os import json import time +import glob import asyncio import logging from typing import Any, Dict @@ -53,10 +58,8 @@ class AioKV: self.lock = asyncio.Lock() async def set(self, key: str, value: Any = None) -> None: - """set a value; if value is omitted or None, stores int(time.time()) - - the timestamp default exists for "mark that i saw/did X at time T" usage. - """ + """set a value; if value is omitted or None, stores int(time.time())""" + self._check_key(key) async with self.lock: cache = await self._load() cache[key] = value if value is not None else int(time.time()) @@ -64,6 +67,7 @@ class AioKV: async def get(self, key: str, default: Any = None) -> Any: """return the value for key, or default if absent""" + self._check_key(key) async with self.lock: cache = await self._load() return cache.get(key, default) @@ -71,6 +75,7 @@ class AioKV: async def delete(self, key: str) -> bool: """delete a key; returns True if removed or absent, False on error""" try: + self._check_key(key) async with self.lock: cache = await self._load() if key in cache: @@ -82,11 +87,7 @@ class AioKV: return False async def clear(self) -> bool: - """remove the backing file entirely; returns True on success, False on error - - a file already absent (or removed concurrently between the check and the remove) - is success — the goal state, no file, is reached. - """ + """remove the backing file entirely; True on success or if already absent, False on error""" try: async with self.lock: try: @@ -103,43 +104,43 @@ class AioKV: async with self.lock: return await self._load() - async def _load(self) -> Dict[str, Any]: - """load the store from disk, returning {} if the file is absent or empty + @staticmethod + def _check_key(key: str) -> None: + """raise ValueError if key is not str (JSON object keys are always strings)""" + if not isinstance(key, str): + raise ValueError(f"aiokv keys must be str, got {type(key).__name__}") - a truncated/corrupt file raises JSONDecodeError, and a file holding valid - JSON that is not an object (e.g. a bare list, number, or null) raises - ValueError — both surfaced to the caller rather than silently masking a - real corruption or returning a non-dict that breaks every other method. - """ - if not await asyncio.to_thread(os.path.exists, self.file): + async def _load(self) -> Dict[str, Any]: + """load the store, returning {} if absent or blank; raises JSONDecodeError on + corrupt content and ValueError on valid-but-non-object JSON rather than masking it""" + try: + async with aiofiles.open(self.file, mode="r", encoding="utf-8") as f: + data = await f.read() + except FileNotFoundError: return {} - async with aiofiles.open(self.file, mode="r", encoding="utf-8") as f: - data = await f.read() - if not data: + if not data.strip(): return {} - loaded = json.loads(data) + loaded = await asyncio.to_thread(json.loads, data) if not isinstance(loaded, dict): raise ValueError(f"store file {self.file} does not hold a JSON object") return loaded async def _save(self, cache: Dict[str, Any]) -> None: - """write the store atomically: temp file in the same dir, then os.replace - - os.replace is atomic on POSIX, so a reader never sees a partial file and a - process crash mid-write leaves the previous good file intact. note this is - process-crash safety, NOT power-loss durability — there is no fsync of the temp - file or the directory, so an OS/power failure could still lose the most recent - write (acceptable here: this is reconstructible single-process state, not a db). - """ - directory = os.path.dirname(self.file) or "." + """write atomically: temp file in the same dir, then os.replace over the realpath'd + target (symlink-safe). process-crash safe (a crash mid-write leaves the prior good + file intact), NOT power-loss safe — no fsync, so an OS/power failure can still lose + the last write (acceptable: reconstructible single-process state, not a db)""" + target = await asyncio.to_thread(os.path.realpath, self.file) + directory = os.path.dirname(target) or "." await asyncio.to_thread(os.makedirs, directory, exist_ok=True) + await self._sweep_stale_tmp(directory) - payload = json.dumps(cache) - tmp = f"{self.file}.{os.getpid()}.tmp" + payload = await asyncio.to_thread(json.dumps, cache, allow_nan=False) + tmp = f"{target}.{os.getpid()}.tmp" try: async with aiofiles.open(tmp, mode="w", encoding="utf-8") as f: await f.write(payload) - await asyncio.to_thread(os.replace, tmp, self.file) + await asyncio.to_thread(os.replace, tmp, target) except Exception: if await asyncio.to_thread(os.path.exists, tmp): try: @@ -148,6 +149,37 @@ class AioKV: log.exception("aiokv: failed to clean up temp file %s", tmp) raise + async def _sweep_stale_tmp(self, directory: str) -> None: + """remove orphaned ..tmp files left by a hard crash of a different, dead process""" + base = os.path.basename(await asyncio.to_thread(os.path.realpath, self.file)) + pattern = os.path.join(directory, f"{base}.*.tmp") + for path in await asyncio.to_thread(glob.glob, pattern): + try: + pid = int(path.rsplit(".", 2)[-2]) + except (ValueError, IndexError): + continue + if pid == os.getpid(): + continue + if await asyncio.to_thread(self._pid_alive, pid): + continue + try: + await asyncio.to_thread(os.remove, path) + except FileNotFoundError: + pass + except Exception: + log.exception("aiokv: failed to sweep stale temp file %s", path) + + @staticmethod + def _pid_alive(pid: int) -> bool: + """check whether pid refers to a live process, without permission to signal counting as alive""" + try: + os.kill(pid, 0) + except ProcessLookupError: + return False + except PermissionError: + return True + return True + # back-compat: this lib was originally named aiocache; legacy call sites using # `aiocache(...)` keep working via this alias. prefer AioKV in new code.