fix: lock the session before opening it, and stop a failed export reading as finished

The session lock was taken after the thing it guards. It lived inside
connected_client, so connect() and _login() had already written the
shared SQLite session by the time the lock existed. Two runs on the
default session each passed their own per-root lock, both opened the same
database, and the loser exited 3 *after* causing the corruption its
message described. The lock now wraps the whole client lifetime, and the
committable-path check runs first so a refused run leaves no lock file
next to a session it was never allowed to create.

--dry-run is inside that lock too. It writes nothing to the export tree,
but it opens the same session file, which is the resource the lock is
about - so running it alongside an export now needs its own --session.

Refusing to mark an export complete required that *nothing* had
succeeded: `failed and not (downloaded or skipped)`. A single
already-present file made skipped non-zero and disabled the guard
outright, so any resume across a partly-complete export could fail every
remaining file and still stamp completed_at with a cursor at
end-of-history. A broken export then answered "did my export finish?"
with a confident yes. The comparison is now against downloaded + skipped;
a legitimate tail of present files still outnumbers its own stray
failures and completes normally.

Also:

- Renewing an expired file reference counted as a failed attempt, so
  expiry on the final attempt burned the last slot and the fresh
  reference was never fetched - reported as "exhausted 3 attempts" after
  two. The refreshed latch already bounds that arm.
- Sidecar repair scanned a fixed window back from EOF for the last
  newline. A partial record larger than the window contains none, so the
  file was truncated to end-window: still unreadable, one megabyte
  shorter. The window now grows until a newline is found.
- --reset-state made the zeroed cursor durable before rotating the old
  sidecar, leaving exactly the mixed-generation log that rotating exists
  to prevent. Rotation now precedes State.open, which is sound only on
  this path because there is no compatibility check to fail.
- Two resets inside one second silently clobbered the first archive
  through os.replace, and the rename was the one here not fsynced.
- title.txt was the only untrusted string written raw. The export root is
  safe because the title never becomes a path component, but cat title.txt
  handed ANSI escapes and a right-to-left override to the operator. It is
  stripped of the same Unicode categories filenames are, from one shared
  set so the two rules cannot drift, and written atomically.

Cost, on the most common path of every resume:

- The post directory was fsynced every post, including posts where every
  file was already present and nothing had been renamed - one fsync per
  post to re-record a directory entry an earlier run had already made
  durable. Gated on an actual download.
- An already-present file was stat'd three times: exists(), stat(), and
  again inside _result. One stat now, reused as the recorded size.

Moving fsync off the loop is not available: the AST scan forbids
to_thread and run_in_executor, deliberately, because parallelism here
buys nothing and escalates flood waits.

Each fix has a test that fails without it, confirmed by reverting the
fix and re-running. 204 tests to 214, coverage unchanged at 93%.
This commit is contained in:
tiennm99 committed 2026-08-22 23:20:00 +07:00
1 parent 3a51f33068
commit c37598e008
8 files changed
+308 -45

No files matched your search

+3
View File
@@ -10,3 +10,6 @@ exports/
__pycache__/ __pycache__/
*.egg-info/ *.egg-info/
.env .env
.coverage
.pytest_cache/
.claude/agent-memory/
+44 -25
View File
@@ -154,22 +154,35 @@ async def _run(args) -> int:
session_path = with_session_suffix( session_path = with_session_suffix(
Path(args.session).expanduser() if args.session else default_session_path()) Path(args.session).expanduser() if args.session else default_session_path())
async with connected_client(session_path, api_id, api_hash) as client: # The session lock guards the SQLite session file, so it has to be held
entity, peer_id = await resolve_entity(client, args.group) # before anything opens that file - connect() and _login() both write to it.
title = (getattr(entity, "title", None) # Taken deeper in, it locked *after* the damage: two runs sharing the default
or getattr(entity, "username", None) or str(peer_id)) # session each passed their own per-root lock, both opened the same database,
root = paths.export_root(Path(args.out).expanduser(), peer_id) # and the loser reported the conflict having already caused it.
assert_not_committable(root, is_dir=True) #
log.info("group %r (id %s) -> %s", title, peer_id, root) # --dry-run is inside the lock too. It writes nothing to the export tree, but
# it opens the same session file, which is the resource this lock is about -
# so a dry-run alongside a real export needs its own --session.
#
# Refuse a committable path before the lock, so a rejected run leaves no
# lock file behind next to a session it was never allowed to create.
assert_not_committable(session_path)
with exclusive(session_path.with_name(session_path.name + ".lock")):
async with connected_client(session_path, api_id, api_hash) as client:
entity, peer_id = await resolve_entity(client, args.group)
title = (getattr(entity, "title", None)
or getattr(entity, "username", None) or str(peer_id))
root = paths.export_root(Path(args.out).expanduser(), peer_id)
assert_not_committable(root, is_dir=True)
log.info("group %r (id %s) -> %s", title, peer_id, root)
if args.dry_run: if args.dry_run:
return await _dry_run(client, entity, root=root, title=title, return await _dry_run(client, entity, root=root, title=title,
peer_id=peer_id, filters=filters, peer_id=peer_id, filters=filters,
min_free=min_free, args=args) min_free=min_free, args=args)
return await _real_run(client, entity, root=root, title=title, return await _real_run(client, entity, root=root, title=title,
peer_id=peer_id, filters=filters, peer_id=peer_id, filters=filters,
min_free=min_free, session_path=session_path, min_free=min_free, args=args)
args=args)
async def _dry_run(client, entity, *, root: Path, title: str, peer_id: int, async def _dry_run(client, entity, *, root: Path, title: str, peer_id: int,
@@ -185,24 +198,30 @@ async def _dry_run(client, entity, *, root: Path, title: str, peer_id: int,
async def _real_run(client, entity, *, root: Path, title: str, peer_id: int, async def _real_run(client, entity, *, root: Path, title: str, peer_id: int,
filters, min_free: int, session_path: Path, args) -> int: filters, min_free: int, args) -> int:
root.mkdir(parents=True, exist_ok=True) root.mkdir(parents=True, exist_ok=True)
session_lock = session_path.with_name(session_path.name + ".lock")
# Both locks before the .part sweep: the sweep is the destructive operation, # The export lock before the .part sweep: the sweep is the destructive
# so locking after it would lock after the damage. Per-root rather than # operation, so locking after it would lock after the damage. Per-root rather
# global, because exporting two different groups concurrently is legitimate. # than global, because exporting two different groups concurrently is
with exclusive(root / LOCK_NAME), exclusive(session_lock): # legitimate - each with its own --session, which the session lock in _run
# now enforces rather than merely documents.
with exclusive(root / LOCK_NAME):
sweep_part_files(root) sweep_part_files(root)
write_title(root, title) write_title(root, title)
# State first: it refuses a filter or chat mismatch (exit 7) before the # Rotate before the cursor is zeroed, and only under --reset-state -
# sidecar is touched. # which is the one path with no compatibility check to fail, since the
state = State.open(root, chat_id=peer_id, chat_title=title, # stored filters are being discarded. Opening state first made the zeroed
filters=filters, reset=args.reset_state) # cursor durable while the previous regime's sidecar was still in place,
# which is exactly the mixed-generation log that rotating exists to
# prevent. Without --reset-state nothing here is touched, so State.open
# still refuses a filter or chat mismatch (exit 7) before any write.
sidecar = Sidecar(root) sidecar = Sidecar(root)
if args.reset_state: if args.reset_state:
sidecar.rotate() sidecar.rotate()
state = State.open(root, chat_id=peer_id, chat_title=title,
filters=filters, reset=args.reset_state)
with sidecar: with sidecar:
posts = iter_posts(client, entity, after_id=state.cursor_id, posts = iter_posts(client, entity, after_id=state.cursor_id,
+65 -17
View File
@@ -280,9 +280,16 @@ async def download_one(fetcher, post_dir_path: Path, msg, *, cfg: Config) -> Res
log.error("cannot build a safe path for message %s: %s - skipping", msg.id, e) log.error("cannot build a safe path for message %s: %s - skipping", msg.id, e)
return _result(msg, FAILED, error=f"unsafe filename: {e}") return _result(msg, FAILED, error=f"unsafe filename: {e}")
if target.exists(): # One stat, not three. An already-present file is the common case on every
if target.stat().st_size > 0: # resume, and reusing this size as the recorded size keeps the whole skip
return _result(msg, SKIPPED, target) # path to a single syscall.
try:
present = target.stat().st_size
except FileNotFoundError:
present = None
if present is not None:
if present > 0:
return _result(msg, SKIPPED, target, size=present)
# Zero bytes means ENOSPC junk, never a complete download - re-fetch. # Zero bytes means ENOSPC junk, never a complete download - re-fetch.
target.unlink() target.unlink()
@@ -291,7 +298,9 @@ async def download_one(fetcher, post_dir_path: Path, msg, *, cfg: Config) -> Res
tmp = target.with_name(target.name + PART_SUFFIX) tmp = target.with_name(target.name + PART_SUFFIX)
refreshed = False refreshed = False
for attempt in range(1, MAX_ATTEMPTS + 1): attempt = 0
while attempt < MAX_ATTEMPTS:
attempt += 1
try: try:
with open(tmp, "wb") as fh: with open(tmp, "wb") as fh:
await fetcher.fetch(msg, fh) await fetcher.fetch(msg, fh)
@@ -315,6 +324,12 @@ async def download_one(fetcher, post_dir_path: Path, msg, *, cfg: Config) -> Res
refreshed = True refreshed = True
if fresh is None or not has_media(fresh): if fresh is None or not has_media(fresh):
return _result(msg, FAILED, error="message deleted") return _result(msg, FAILED, error="message deleted")
# Renewing a reference is not a failed attempt. Counting it burned
# the last slot when expiry landed on the final attempt, so the
# fresh reference was never actually fetched and the message was
# reported as "exhausted N attempts" having been tried N-1 times.
# The `refreshed` latch above bounds this arm, so it cannot loop.
attempt -= 1
msg = fresh msg = fresh
continue continue
except _UNAVAILABLE_ERRORS as e: except _UNAVAILABLE_ERRORS as e:
@@ -369,13 +384,15 @@ async def download_one(fetcher, post_dir_path: Path, msg, *, cfg: Config) -> Res
def _result(msg, status: str, target: Path | None = None, def _result(msg, status: str, target: Path | None = None,
error: str | None = None) -> Result: error: str | None = None, size: int | None = None) -> Result:
f = getattr(msg, "file", None) f = getattr(msg, "file", None)
if size is None and target is not None:
size = target.stat().st_size
return Result( return Result(
message_id=msg.id, message_id=msg.id,
status=status, status=status,
path=target, path=target,
size=target.stat().st_size if target is not None else None, size=size,
declared_size=media_size(msg), declared_size=media_size(msg),
name=getattr(f, "name", None) or f"{media_kind(msg)}{file_ext(msg)}", name=getattr(f, "name", None) or f"{media_kind(msg)}{file_ext(msg)}",
mime=getattr(f, "mime_type", None), mime=getattr(f, "mime_type", None),
@@ -416,7 +433,11 @@ async def run_download(posts, *, fetcher, state, sidecar, root: Path,
result = await download_one(fetcher, target_dir, msg, cfg=cfg) result = await download_one(fetcher, target_dir, msg, cfg=cfg)
results.append(result) results.append(result)
_tally(totals, result, post.post_id) _tally(totals, result, post.post_id)
if target_dir.exists(): # Only a rename needs its directory entry forced. When every file in
# the post was already present nothing was renamed, so a resume
# across a mostly-complete export no longer pays one fsync per post
# to make a directory entry durable that a previous run already did.
if any(r.status == DOWNLOADED for r in results) and target_dir.exists():
fsync_dir(target_dir) fsync_dir(target_dir)
# A post may be empty (everything filtered out) or text-only under # A post may be empty (everything filtered out) or text-only under
@@ -445,15 +466,25 @@ async def run_download(posts, *, fetcher, state, sidecar, root: Path,
# at the last post that happened to contain a file - and the next run would # at the last post that happened to contain a file - and the next run would
# re-sweep that whole tail to download nothing. # re-sweep that whole tail to download nothing.
state.commit(last_complete) state.commit(last_complete)
if totals.failed and not (totals.downloaded or totals.skipped): if totals.failed > totals.downloaded + totals.skipped:
# Every single file failed and none succeeded. That is what a dead # Failures outnumber everything that worked. That is what a dead media DC
# session looks like from in here - ConnectionError is retryable, so the # looks like from in here - ConnectionError is retryable, so the loop
# loop walks the whole history failing everything and returns normally. # walks the whole history failing everything and returns normally.
# Setting completed_at on that would answer "did my export finish?" with # Setting completed_at on that would answer "did my export finish?" with
# a confident yes over an empty tree. The cursor still stands: the posts # a confident yes over a mostly-empty tree.
# were handled, and their failures are in the sidecar. #
log.error("every file failed (%d) and none succeeded - not marking this " # Compared against downloaded + skipped, not against zero: requiring
"export complete; check connectivity and re-run", totals.failed) # *nothing* to have succeeded meant a single already-present file
# disabled the guard entirely, so any resume over a partly-complete
# export could fail every remaining file and still be marked finished.
# A legitimate tail of already-present files still outnumbers its own
# stray failures, so it completes normally.
#
# The cursor still stands either way: the posts were handled, and their
# failures are in the sidecar.
log.error("%d files failed, outnumbering the %d downloaded and %d already "
"present - not marking this export complete; check connectivity "
"and re-run", totals.failed, totals.downloaded, totals.skipped)
else: else:
state.mark_completed() state.mark_completed()
return totals return totals
@@ -474,6 +505,23 @@ def _tally(totals: Totals, result: Result, post_id: int) -> None:
def write_title(root: Path, title: str) -> None: def write_title(root: Path, title: str) -> None:
"""The human-readable group name lives here, as data - never as a path """The human-readable group name lives here, as data - never as a path
component. Written once; a later rename updates it without moving anything.""" component. Written once; a later rename updates it without moving anything.
Server-supplied and editable by any admin, so it is stripped of the
unprintable categories before landing on disk: the export root is protected
by never putting the title in a path at all, but `cat title.txt` would
otherwise feed ANSI escapes or a right-to-left override to the operator's
terminal.
Atomic like every other write here. A torn title.txt is only cosmetic, but it
was the one file in the tree that could be observed half-written.
"""
root.mkdir(parents=True, exist_ok=True) root.mkdir(parents=True, exist_ok=True)
(root / TITLE_NAME).write_text((title or "") + "\n", encoding="utf-8") target = root / TITLE_NAME
tmp = target.with_name(target.name + PART_SUFFIX)
with open(tmp, "w", encoding="utf-8") as f:
f.write(paths.strip_unprintable(title or "") + "\n")
f.flush()
os.fsync(f.fileno())
os.replace(tmp, target)
fsync_dir(root)
+14
View File
@@ -56,6 +56,20 @@ MAX_NAME_BYTES = 200
MAX_EXT_BYTES = 16 MAX_EXT_BYTES = 16
def strip_unprintable(text: str) -> str:
"""Strip the categories that make a string unsafe to *print*, leaving it
otherwise intact.
sanitize() below reduces an untrusted string to a single path component.
This does the much smaller job that untrusted text written as *data* needs -
title.txt - where spaces, punctuation and slashes are all legitimate content,
but an ANSI escape or a U+202E override still reaches the operator's terminal
the moment they `cat` the file. Same category set, so the two rules cannot
drift apart.
"""
return "".join(ch for ch in str(text) if not _is_disallowed(ch))
def export_root(out_dir: Path, peer_id: int) -> Path: def export_root(out_dir: Path, peer_id: int) -> Path:
"""`<out>/g<chat_id>` - the directory name is the chat id, never the title. """`<out>/g<chat_id>` - the directory name is the chat id, never the title.
+21 -3
View File
@@ -14,6 +14,8 @@ import os
from datetime import UTC, datetime from datetime import UTC, datetime
from pathlib import Path from pathlib import Path
from .state import fsync_dir
log = logging.getLogger("telegram_exporter.sidecar") log = logging.getLogger("telegram_exporter.sidecar")
SIDECAR_NAME = "messages.jsonl" SIDECAR_NAME = "messages.jsonl"
@@ -58,10 +60,18 @@ class Sidecar:
with open(self.path, "rb+") as f: with open(self.path, "rb+") as f:
f.seek(0, os.SEEK_END) f.seek(0, os.SEEK_END)
end = f.tell() end = f.tell()
# Grow the window until a newline is actually found. A record longer
# than the window contains no newline at all, and truncating to
# `end - window` would then cut in the middle of an earlier, valid
# record - leaving the file just as unreadable, one megabyte shorter.
window = min(end, 1 << 20) window = min(end, 1 << 20)
f.seek(end - window) while True:
tail = f.read(window) f.seek(end - window)
cut = tail.rfind(b"\n") tail = f.read(window)
cut = tail.rfind(b"\n")
if cut != -1 or window == end:
break
window = min(end, window * 2)
trailing = tail[cut + 1:] trailing = tail[cut + 1:]
if not trailing: if not trailing:
return # file ends on a newline: intact return # file ends on a newline: intact
@@ -83,7 +93,15 @@ class Sidecar:
return None return None
stamp = datetime.now(UTC).strftime("%Y%m%dT%H%M%SZ") stamp = datetime.now(UTC).strftime("%Y%m%dT%H%M%SZ")
target = self.root / f"messages-{stamp}.jsonl" target = self.root / f"messages-{stamp}.jsonl"
# Two resets inside the same second must not clobber the first archive:
# rotating instead of deleting is only worth doing if the old log stays
# inspectable, and os.replace would have overwritten it silently.
serial = 1
while target.exists():
target = self.root / f"messages-{stamp}-{serial}.jsonl"
serial += 1
os.replace(self.path, target) os.replace(self.path, target)
fsync_dir(self.root) # a rename is not durable until its dir entry is
log.info("rotated %s -> %s", self.path.name, target.name) log.info("rotated %s -> %s", self.path.name, target.name)
return target return target
+32
View File
@@ -205,3 +205,35 @@ def test_a_second_concurrent_run_exits_three(wired, monkeypatch, tmp_path):
with cli.exclusive(root / cli.LOCK_NAME): with cli.exclusive(root / cli.LOCK_NAME):
code = run_cli(monkeypatch, "--group", "@mygroup", "--out", str(out)) code = run_cli(monkeypatch, "--group", "@mygroup", "--out", str(out))
assert code == 3 assert code == 3
def test_a_session_lock_conflict_is_refused_before_the_session_is_opened(
wired, monkeypatch, tmp_path):
# The session lock was taken after connect() and _login() had already written
# the shared SQLite session, so the loser reported the conflict having caused
# it. The session file must not exist at all when the lock is contended.
session = tmp_path / "state" / "tg-export" / "default.session"
session.parent.mkdir(parents=True, exist_ok=True)
with cli.exclusive(session.with_name(session.name + ".lock")):
code = run_cli(monkeypatch, "--group", "@mygroup",
"--out", str(tmp_path / "exports"))
assert code == 3
assert not session.exists()
assert wired.downloads == []
def test_dry_run_contends_for_the_session_lock_too(wired, monkeypatch, tmp_path):
# --dry-run writes nothing to the export tree, but it opens the same session
# file - which is the resource this lock is about - so it needs its own
# --session to run alongside a real export.
session = tmp_path / "state" / "tg-export" / "default.session"
session.parent.mkdir(parents=True, exist_ok=True)
with cli.exclusive(session.with_name(session.name + ".lock")):
code = run_cli(monkeypatch, "--group", "@mygroup", "--dry-run",
"--out", str(tmp_path / "exports"))
assert code == 3
assert not session.exists()
+97
View File
@@ -35,6 +35,7 @@ from telegram_exporter.downloader import (
download_one, download_one,
run_download, run_download,
sweep_part_files, sweep_part_files,
TITLE_NAME,
write_title, write_title,
) )
from telegram_exporter.session import Abort from telegram_exporter.session import Abort
@@ -525,3 +526,99 @@ def test_limit_zero_connects_and_does_nothing(tmp_path):
assert totals.posts == 0 assert totals.posts == 0
assert state.cursor_id == 0 assert state.cursor_id == 0
assert not (root / "1042").exists() assert not (root / "1042").exists()
def test_a_reference_renewed_on_the_final_attempt_is_actually_fetched(tmp_path):
# Renewing a file reference is not a failed attempt. Counting it consumed
# the last slot, so expiry landing on the final attempt reported "exhausted
# 3 attempts" for a message that was only ever tried twice - and the docstring
# on that arm says a 20-hour run *will* reach it.
root, target_dir = post_dir_for(tmp_path)
msg = FakeMsg(1042, kind="photo", size=len(PAYLOAD))
fetcher = FakeFetcher(raises=[ConnectionResetError(), ConnectionResetError(),
FileReferenceExpiredError(request=None)])
result = run(download_one(fetcher, target_dir, msg, cfg=cfg(tmp_path)))
assert result.status == DOWNLOADED
assert fetcher.refresh_calls == 1
assert (target_dir / "1042_photo.jpg").read_bytes() == PAYLOAD
def test_an_already_present_file_reports_its_size_without_restatting(tmp_path):
root, target_dir = post_dir_for(tmp_path)
(target_dir / "1042_photo.jpg").write_bytes(b"abcdef")
msg = FakeMsg(1042, kind="photo", size=len(PAYLOAD))
fetcher = FakeFetcher()
result = run(download_one(fetcher, target_dir, msg, cfg=cfg(tmp_path)))
assert result.status == SKIPPED
assert result.size == 6
assert fetcher.calls == 0
def test_a_post_whose_files_are_all_present_does_not_fsync_its_directory(
tmp_path, monkeypatch):
# The post directory is fsynced to make a rename durable. When every file was
# already on disk nothing was renamed, so a resume across a mostly-complete
# export should not pay one fsync per post for a directory entry an earlier
# run already made durable.
calls = []
monkeypatch.setattr(downloader, "fsync_dir", lambda p: calls.append(p))
root = paths.export_root(tmp_path, -1001)
for message_id in (1042, 1043):
path = paths.post_dir(root, 1042) / f"{message_id}_photo.jpg"
path.parent.mkdir(parents=True, exist_ok=True)
path.write_bytes(PAYLOAD)
_, _, totals, _ = drive(tmp_path, [album(1042, [1042, 1043])])
assert totals.skipped == 2 and totals.downloaded == 0
assert calls == []
def test_a_post_with_one_new_file_still_fsyncs_its_directory(tmp_path, monkeypatch):
calls = []
monkeypatch.setattr(downloader, "fsync_dir", lambda p: calls.append(p))
_, _, totals, _ = drive(tmp_path, [album(1042, [1042, 1043])])
assert totals.downloaded == 2
assert calls == [paths.post_dir(paths.export_root(tmp_path, -1001), 1042)]
def test_a_resume_where_failures_outnumber_successes_is_not_marked_complete(tmp_path):
# One already-present file used to disable this guard outright: it required
# that *nothing* had succeeded, so any resume across a partly-complete export
# could fail every remaining file and still record the export as finished.
class MediaDcIsDead(FakeFetcher):
async def fetch(self, msg, fh):
self.calls += 1
raise ConnectionError("Cannot send requests while disconnected")
root = paths.export_root(tmp_path, -1001)
present = paths.post_dir(root, 1) / "1_photo.jpg"
present.parent.mkdir(parents=True, exist_ok=True)
present.write_bytes(PAYLOAD)
root, state, totals, _ = drive(
tmp_path, [album(1, [1]), album(2, [2]), album(3, [3])],
fetcher=MediaDcIsDead())
assert totals.skipped == 1 and totals.failed == 2 and totals.downloaded == 0
assert State.read(root)["completed_at"] is None
assert state.cursor_id == 3 # the cursor still stands
def test_the_group_title_is_written_without_terminal_escapes(tmp_path):
# Server-supplied and admin-editable. The export root is safe because the
# title never becomes a path component, but `cat title.txt` would still have
# handed ANSI escapes and a right-to-left override to the operator.
root = paths.export_root(tmp_path, -1001)
write_title(root, "My \x1b[31mGroup\x1b[0m ‮gpj.exe")
text = (root / "title.txt").read_text()
assert "\x1b" not in text and "‮" not in text
assert "Group" in text and "gpj.exe" in text
assert not (root / (TITLE_NAME + ".part")).exists() # atomic, nothing left over
+32
View File
@@ -119,3 +119,35 @@ def test_sender_name_is_only_taken_from_an_already_cached_sender(tmp_path):
assert records[0]["sender_name"] == "Alice" assert records[0]["sender_name"] == "Alice"
assert records[1]["sender_name"] is None assert records[1]["sender_name"] is None
assert records[1]["sender_id"] == 778 assert records[1]["sender_id"] == 778
def test_a_partial_record_larger_than_the_scan_window_keeps_earlier_records(tmp_path):
# Scanning back a fixed window for the last newline finds none when the
# partial record is bigger than the window. Truncating to `end - window`
# then cuts inside the partial record and leaves the file just as unreadable,
# one window shorter - so the window has to grow until a newline is found.
path = tmp_path / SIDECAR_NAME
good = '{"message_id": 1}'
partial = '{"message_id": 2, "caption": "' + "x" * (2 << 20)
path.write_text(good + "\n" + partial)
with Sidecar(tmp_path):
pass
assert lines(tmp_path) == [good]
def test_two_rotations_in_the_same_second_both_stay_inspectable(tmp_path):
# os.replace onto a one-second stamp silently overwrote the first archive,
# which is the whole reason --reset-state rotates instead of deleting.
path = tmp_path / SIDECAR_NAME
path.write_text('{"message_id": 1}\n')
first = Sidecar(tmp_path).rotate()
path.write_text('{"message_id": 2}\n')
second = Sidecar(tmp_path).rotate()
assert first is not None and second is not None
assert first != second
assert first.exists() and second.exists()
assert json.loads(first.read_text())["message_id"] == 1
assert json.loads(second.read_text())["message_id"] == 2