Skip to content

Commit 5466132

Browse files
committed
Simplify runtime journal recovery
1 parent a43c2e2 commit 5466132

2 files changed

Lines changed: 13 additions & 31 deletions

File tree

‎src/_pytask/journal.py‎

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -28,7 +28,7 @@ def append(self, payload: msgspec.Struct) -> None:
2828
journal_file.write(msgspec.json.encode(payload) + b"\n")
2929

3030
def read(self) -> list[T]:
31-
"""Read entries, stopping at the first invalid line."""
31+
"""Read entries, clearing the journal on decode errors."""
3232
if not self.path.exists():
3333
return []
3434

@@ -39,7 +39,8 @@ def read(self) -> list[T]:
3939
try:
4040
entries.append(msgspec.json.decode(line, type=self.type_))
4141
except msgspec.DecodeError:
42-
break
42+
self.delete()
43+
return []
4344
return entries
4445

4546
def delete(self) -> None:

‎src/_pytask/runtime_store.py‎

Lines changed: 10 additions & 29 deletions
Original file line numberDiff line numberDiff line change
@@ -7,7 +7,6 @@
77
from typing import TYPE_CHECKING
88

99
import msgspec
10-
from packaging.version import Version
1110

1211
from _pytask.journal import JsonlJournal
1312

@@ -19,14 +18,6 @@
1918
CURRENT_RUNTIME_VERSION = "1"
2019

2120

22-
class RuntimeStoreError(Exception):
23-
"""Raised when reading or writing runtime files fails."""
24-
25-
26-
class RuntimeStoreVersionError(RuntimeStoreError):
27-
"""Raised when a runtime file version is not supported."""
28-
29-
3021
class _RuntimeEntry(msgspec.Struct):
3122
id: str
3223
date: float
@@ -63,15 +54,12 @@ def _read_runtimes(path: Path) -> _RuntimeFile | None:
6354
try:
6455
data = msgspec.json.decode(path.read_bytes(), type=_RuntimeFile)
6556
except msgspec.DecodeError:
66-
msg = "Runtime file has invalid format."
67-
raise RuntimeStoreError(msg) from None
57+
path.unlink()
58+
return None
6859

69-
if Version(data.runtime_version) != Version(CURRENT_RUNTIME_VERSION):
70-
msg = (
71-
f"Unsupported runtime-version {data.runtime_version!r}. "
72-
f"Current version is {CURRENT_RUNTIME_VERSION}."
73-
)
74-
raise RuntimeStoreVersionError(msg)
60+
if data.runtime_version != CURRENT_RUNTIME_VERSION:
61+
path.unlink()
62+
return None
7563
return data
7664

7765

@@ -87,12 +75,9 @@ def _read_journal(
8775
) -> list[_RuntimeJournalEntry]:
8876
entries = journal.read()
8977
for entry in entries:
90-
if Version(entry.runtime_version) != Version(CURRENT_RUNTIME_VERSION):
91-
msg = (
92-
f"Unsupported runtime-version {entry.runtime_version!r}. "
93-
f"Current version is {CURRENT_RUNTIME_VERSION}."
94-
)
95-
raise RuntimeStoreVersionError(msg)
78+
if entry.runtime_version != CURRENT_RUNTIME_VERSION:
79+
journal.delete()
80+
return []
9681
return entries
9782

9883

@@ -147,12 +132,8 @@ def from_root(cls, root: Path) -> RuntimeState:
147132
def _rebuild_index(self) -> None:
148133
self._index = {entry.id: entry for entry in self.runtimes.task}
149134

150-
@staticmethod
151-
def _task_id(task: PTask) -> str:
152-
return task.name
153-
154135
def update_task(self, task: PTask, start: float, end: float) -> None:
155-
task_id = self._task_id(task)
136+
task_id = task.name
156137
entry = _RuntimeEntry(id=task_id, date=start, duration=end - start)
157138
self._index[entry.id] = entry
158139
self.runtimes = _RuntimeFile(
@@ -170,7 +151,7 @@ def update_task(self, task: PTask, start: float, end: float) -> None:
170151
self._dirty = True
171152

172153
def get_duration(self, task: PTask) -> float | None:
173-
task_id = self._task_id(task)
154+
task_id = task.name
174155
entry = self._index.get(task_id)
175156
if entry is None:
176157
return None

0 commit comments

Comments
 (0)