EvaluationWorkspace
This commit is contained in:
@@ -0,0 +1,168 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import os
|
||||
import stat
|
||||
import tempfile
|
||||
import time
|
||||
from collections.abc import Callable
|
||||
from contextlib import contextmanager
|
||||
from pathlib import Path
|
||||
from typing import Any, BinaryIO, TypeVar, cast
|
||||
|
||||
JsonValue = dict[str, Any] | list[Any] | str | int | float | bool | None
|
||||
T = TypeVar("T", bound=JsonValue)
|
||||
_MISSING = object()
|
||||
|
||||
|
||||
class JsonLockTimeout(TimeoutError):
|
||||
pass
|
||||
|
||||
|
||||
def read_json(path: str | Path, *, default: T | object = _MISSING) -> T:
|
||||
target = Path(path)
|
||||
try:
|
||||
with target.open("r", encoding="utf-8") as stream:
|
||||
return cast(T, json.load(stream))
|
||||
except FileNotFoundError:
|
||||
if default is _MISSING:
|
||||
raise
|
||||
return cast(T, default)
|
||||
|
||||
|
||||
def _sync_directory(directory: Path) -> None:
|
||||
if os.name == "nt" or not hasattr(os, "O_DIRECTORY"):
|
||||
return
|
||||
descriptor = os.open(directory, os.O_RDONLY | os.O_DIRECTORY)
|
||||
try:
|
||||
os.fsync(descriptor)
|
||||
finally:
|
||||
os.close(descriptor)
|
||||
|
||||
|
||||
def _atomic_write(path: Path, payload: bytes) -> None:
|
||||
path = path.expanduser()
|
||||
path.parent.mkdir(parents=True, exist_ok=True)
|
||||
previous_mode = None
|
||||
try:
|
||||
previous_mode = stat.S_IMODE(path.stat().st_mode)
|
||||
except FileNotFoundError:
|
||||
pass
|
||||
|
||||
descriptor, temporary_name = tempfile.mkstemp(
|
||||
dir=path.parent,
|
||||
prefix=f".{path.name}.",
|
||||
suffix=".tmp",
|
||||
)
|
||||
temporary = Path(temporary_name)
|
||||
try:
|
||||
with os.fdopen(descriptor, "wb") as stream:
|
||||
stream.write(payload)
|
||||
stream.flush()
|
||||
os.fsync(stream.fileno())
|
||||
if previous_mode is not None:
|
||||
os.chmod(temporary, previous_mode)
|
||||
os.replace(temporary, path)
|
||||
_sync_directory(path.parent)
|
||||
finally:
|
||||
if temporary.exists():
|
||||
temporary.unlink()
|
||||
|
||||
|
||||
def atomic_write_text(
|
||||
path: str | Path,
|
||||
text: str,
|
||||
*,
|
||||
encoding: str = "utf-8",
|
||||
) -> None:
|
||||
_atomic_write(Path(path), text.encode(encoding))
|
||||
|
||||
|
||||
def atomic_write_json(
|
||||
path: str | Path,
|
||||
value: JsonValue,
|
||||
*,
|
||||
indent: int | None = 2,
|
||||
ensure_ascii: bool = False,
|
||||
sort_keys: bool = False,
|
||||
) -> None:
|
||||
serialized = json.dumps(
|
||||
value,
|
||||
indent=indent,
|
||||
ensure_ascii=ensure_ascii,
|
||||
sort_keys=sort_keys,
|
||||
)
|
||||
atomic_write_text(path, serialized + "\n")
|
||||
|
||||
|
||||
def _prepare_lock_file(stream: BinaryIO) -> None:
|
||||
stream.seek(0, os.SEEK_END)
|
||||
if stream.tell() == 0:
|
||||
stream.write(b"\0")
|
||||
stream.flush()
|
||||
|
||||
|
||||
def _try_lock(stream: BinaryIO) -> None:
|
||||
stream.seek(0)
|
||||
if os.name == "nt":
|
||||
import msvcrt
|
||||
|
||||
msvcrt.locking(stream.fileno(), msvcrt.LK_NBLCK, 1)
|
||||
else:
|
||||
import fcntl
|
||||
|
||||
fcntl.flock(stream.fileno(), fcntl.LOCK_EX | fcntl.LOCK_NB)
|
||||
|
||||
|
||||
def _unlock(stream: BinaryIO) -> None:
|
||||
stream.seek(0)
|
||||
if os.name == "nt":
|
||||
import msvcrt
|
||||
|
||||
msvcrt.locking(stream.fileno(), msvcrt.LK_UNLCK, 1)
|
||||
else:
|
||||
import fcntl
|
||||
|
||||
fcntl.flock(stream.fileno(), fcntl.LOCK_UN)
|
||||
|
||||
|
||||
@contextmanager
|
||||
def json_file_lock(path: str | Path, *, timeout: float = 10.0):
|
||||
target = Path(path).expanduser()
|
||||
target.parent.mkdir(parents=True, exist_ok=True)
|
||||
lock_path = target.with_name(f".{target.name}.lock")
|
||||
deadline = time.monotonic() + timeout
|
||||
|
||||
with lock_path.open("a+b") as stream:
|
||||
_prepare_lock_file(stream)
|
||||
while True:
|
||||
try:
|
||||
_try_lock(stream)
|
||||
break
|
||||
except (BlockingIOError, PermissionError, OSError) as exc:
|
||||
if time.monotonic() >= deadline:
|
||||
raise JsonLockTimeout(
|
||||
f"Could not acquire JSON lock within {timeout:.1f}s: {lock_path}"
|
||||
) from exc
|
||||
time.sleep(0.05)
|
||||
try:
|
||||
yield
|
||||
finally:
|
||||
_unlock(stream)
|
||||
|
||||
|
||||
def atomic_update_json(
|
||||
path: str | Path,
|
||||
update: Callable[[T], T | None],
|
||||
*,
|
||||
default_factory: Callable[[], T],
|
||||
timeout: float = 10.0,
|
||||
indent: int | None = 2,
|
||||
) -> T:
|
||||
target = Path(path)
|
||||
with json_file_lock(target, timeout=timeout):
|
||||
current = read_json(target, default=default_factory())
|
||||
updated = update(current)
|
||||
result = current if updated is None else updated
|
||||
atomic_write_json(target, result, indent=indent)
|
||||
return result
|
||||
Reference in New Issue
Block a user