fix(json_helpers): write atomically by default and warn on empty reads

Every caller writes a whole JSON document, and a plain open/write can
leave a torn or empty file when the process dies or another writer
overlaps, which readfile then returns as {}. Writes now go through the
temp file and replace unless atomic=False or the mode is append. A read
waits out an open that Windows refuses while a replace of the same name
is in flight, and an empty file is reported when the read is not
silent, since a zero-byte config.json otherwise resets settings without
a trace.
This commit is contained in:
CalamitousFelicitousness
2026-09-13 19:01:04 +01:00
parent b39639ce2e
commit 68c64fb4bc
2 changed files with 122 additions and 12 deletions
+28 -10
View File
@@ -14,7 +14,7 @@ path_locks_guard = threading.Lock()
def path_lock(filename: str | os.PathLike[str]) -> threading.RLock:
"""One lock per file path; threads of this process serialize on it, other processes are covered by atomic replace."""
key = os.path.normcase(os.path.abspath(filename))
key = os.path.normcase(os.path.realpath(filename))
with path_locks_guard:
lock = path_locks.get(key)
if lock is None:
@@ -24,6 +24,19 @@ def path_lock(filename: str | os.PathLike[str]) -> threading.RLock:
return lock
def read_bytes(filename: str | os.PathLike[str], attempts: int = 5, delay: float = 0.01) -> bytes:
"""Read a whole file, waiting out an open that Windows refuses while a replace of the same name is in flight."""
for attempt in range(attempts):
try:
with open(filename, "rb") as file:
return file.read()
except PermissionError:
if attempt == attempts - 1:
raise
time.sleep(delay)
return b""
def replace_file(source: str, target: str, attempts: int = 10, delay: float = 0.05):
"""os.replace that waits out a target another handle holds open, which Windows reports as a permission error."""
for attempt in range(attempts):
@@ -49,11 +62,12 @@ def readfile(filename: str | os.PathLike[str], silent: bool = False, lock: bool
with path_lock(filename) if lock else contextlib.nullcontext():
try:
t0 = time.time()
with open(filename, "rb") as file:
b = file.read()
if len(b) == 0:
return {} if as_type == "dict" else []
data = orjson.loads(b) # pylint: disable=no-member
b = read_bytes(filename)
if len(b) == 0:
if not silent:
log.warning(f'Read: file="{filename}" empty')
return {} if as_type == "dict" else []
data = orjson.loads(b) # pylint: disable=no-member
t1 = time.time()
if not silent:
fn = f"{sys._getframe(2).f_code.co_name}:{sys._getframe(1).f_code.co_name}" # pylint: disable=protected-access
@@ -81,8 +95,8 @@ def readfile(filename: str | os.PathLike[str], silent: bool = False, lock: bool
return data
def writefile(obj: dict | list, filename: str | os.PathLike[str], mode="w", silent=False, atomic=False):
"""Write obj as JSON; writes to the same path from this process run one at a time, in call order."""
def writefile(obj: dict | list, filename: str | os.PathLike[str], mode="w", silent=False, atomic=True):
"""Write obj as JSON through a temp file and replace; writes to the same path from this process run one at a time, in call order."""
import copy
import tempfile
@@ -90,6 +104,9 @@ def writefile(obj: dict | list, filename: str | os.PathLike[str], mode="w", sile
log.error(f'Save: file="{filename}" not a valid object: {obj}')
return str(obj)
if mode != "w":
atomic = False # append cannot go through a temp file
with path_lock(filename):
try:
t0 = time.time()
@@ -109,13 +126,14 @@ def writefile(obj: dict | list, filename: str | os.PathLike[str], mode="w", sile
try:
if atomic:
fd, temp_name = tempfile.mkstemp(dir=os.path.dirname(os.path.abspath(filename)), prefix=f"{os.path.basename(filename)}.", suffix=".tmp")
target = os.path.realpath(filename) # replace the file a symlink points at, not the symlink
fd, temp_name = tempfile.mkstemp(dir=os.path.dirname(target), prefix=f"{os.path.basename(target)}.", suffix=".tmp")
try:
with os.fdopen(fd, mode, encoding="utf8") as f:
f.write(output)
f.flush()
os.fsync(f.fileno())
replace_file(temp_name, filename)
replace_file(temp_name, target)
except BaseException:
with contextlib.suppress(OSError):
os.remove(temp_name)
+94 -2
View File
@@ -10,6 +10,10 @@ Covers:
- locked readers never see a torn file while unlocked writers rewrite it in place
- a lock file left behind by the former file lock is removed
- concurrent inserts into a shared dict never drop a save
- default writes are atomic, so unlocked readers never see a torn file either
- an atomic write through a symlink replaces the file it points at and keeps the link
- a read waits out an open that is refused while a replace is in flight
- an empty file reads as empty and is reported unless the read is silent
- readfile returns what writefile wrote, as dict and as list
No running server required.
@@ -25,6 +29,7 @@ import shutil
import sys
import tempfile
import threading
import time
script_dir = os.path.dirname(os.path.dirname(os.path.abspath(__file__)))
sys.path.insert(0, script_dir)
@@ -36,17 +41,22 @@ from modules import json_helpers # pylint: disable=wrong-import-position
class ErrorLog:
"""Counts error lines from json_helpers while passing everything else to the real logger."""
"""Counts error and warning lines from json_helpers while passing everything else to the real logger."""
def __init__(self, inner):
self.inner = inner
self.errors = []
self.warnings = []
self.lock = threading.Lock()
def error(self, message, *args, **kwargs):
with self.lock:
self.errors.append(str(message))
def warning(self, message, *args, **kwargs):
with self.lock:
self.warnings.append(str(message))
def count(self, needle):
return sum(1 for m in self.errors if needle in m)
@@ -195,6 +205,85 @@ def test_concurrent_inserts_keep_every_save(folder, threads=4, per_thread=40):
assert len(json.load(f)) == threads * per_thread, 'last save is incomplete'
def test_default_writes_never_tear(folder, writers=4, per_writer=20, readers=2):
jh = fresh_helper()
target = os.path.join(folder, 'default.json')
jh.writefile({'seed': 0}, target, silent=True)
stop = threading.Event()
torn = []
def reader(_n):
while not stop.is_set():
if not jh.readfile(target, silent=True, as_type='dict'):
torn.append(1)
time.sleep(0.001) # real readers do not spin; a spinning reader starves the writer's retry on Windows
def writer(n):
for i in range(per_writer):
jh.writefile({'writer': n, 'i': i, 'blob': 'x' * (40000 + 1000 * n)}, target, silent=True)
pool = [threading.Thread(target=reader, args=(r,)) for r in range(readers)]
for t in pool:
t.start()
run_threads(writer, writers)
stop.set()
for t in pool:
t.join()
assert torn == [], f'{len(torn)} unlocked reads returned nothing'
assert jh.log.errors == [], jh.log.errors
assert leftovers(folder, target) == [], f'temp files left: {leftovers(folder, target)}'
def test_atomic_write_keeps_symlink(folder):
jh = fresh_helper()
real = os.path.join(folder, 'real.json')
link = os.path.join(folder, 'link.json')
jh.writefile({'v': 0}, real, silent=True)
try:
os.symlink(real, link)
except OSError as e:
log.info(f' SKIP symlink not available: {e}')
return
jh.writefile({'v': 1}, link, silent=True)
assert os.path.islink(link), 'symlink was replaced by a file'
assert jh.readfile(real, silent=True, as_type='dict') == {'v': 1}, 'target of the symlink not updated'
assert sorted(os.listdir(folder)) == ['link.json', 'real.json'], os.listdir(folder)
assert jh.log.errors == [], jh.log.errors
def test_read_waits_out_refused_open(folder):
jh = fresh_helper()
target = os.path.join(folder, 'refused.json')
jh.writefile({'x': 1}, target, silent=True)
calls = []
real_open = open
def refuse_twice(*args, **kwargs):
calls.append(1)
if len(calls) <= 2:
raise PermissionError(13, 'replace in flight')
return real_open(*args, **kwargs)
jh.open = refuse_twice # module globals shadow the builtin inside read_bytes
try:
assert jh.readfile(target, silent=True, as_type='dict') == {'x': 1}
finally:
del jh.open
assert len(calls) == 3, f'{len(calls)} open attempts'
assert jh.log.errors == [], jh.log.errors
def test_empty_file_is_reported(folder):
jh = fresh_helper()
target = os.path.join(folder, 'empty.json')
with open(target, 'w', encoding='utf8'):
pass
assert jh.readfile(target, as_type='dict') == {}
assert jh.readfile(target, silent=True, as_type='list') == []
assert len(jh.log.warnings) == 1 and 'empty' in jh.log.warnings[0], jh.log.warnings
assert jh.log.errors == [], jh.log.errors
def test_roundtrip(folder):
jh = fresh_helper()
target = os.path.join(folder, 'roundtrip.json')
@@ -216,6 +305,10 @@ def run_all():
test_locked_readers_never_see_torn_writes,
test_legacy_lock_file_removed,
test_concurrent_inserts_keep_every_save,
test_default_writes_never_tear,
test_atomic_write_keeps_symlink,
test_read_waits_out_refused_open,
test_empty_file_is_reported,
test_roundtrip,
]
passed = 0
@@ -236,7 +329,6 @@ def run_all():
if __name__ == '__main__':
import time
t0 = time.time()
ok = run_all()
log.warning(f'Total time: {time.time() - t0:.2f}s')