-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathrecorder.py
More file actions
71 lines (59 loc) · 2.61 KB
/
Copy pathrecorder.py
File metadata and controls
71 lines (59 loc) · 2.61 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
"""Append-only run recorder.
A `RunRecorder` owns exactly one run directory. Every signal produced
during a run — lessons started/finished, questions answered, per-step
training losses, arbitrary architecture-defined events — funnels through
this object and lands in `events.jsonl` or `training_log.jsonl`.
Design rules:
- `events.jsonl` is the single source of truth. Everything else in the
run directory (`metrics.json`, `curves.json`, `predictions.jsonl`) is
a *derived view* produced at close-time.
- Writes are append-only and flushed on every event so a crashed run
still leaves a partial but readable record.
- `training_log.jsonl` exists as a separate high-frequency channel so
`events.jsonl` stays scannable by eye and by `jq`.
"""
from __future__ import annotations
import json
from datetime import datetime, timezone
from pathlib import Path
from types import TracebackType
from typing import Any
def _now_iso() -> str:
return datetime.now(timezone.utc).isoformat(timespec="milliseconds")
class RunRecorder:
def __init__(self, run_dir: Path) -> None:
self.run_dir = Path(run_dir)
self.run_dir.mkdir(parents=True, exist_ok=True)
(self.run_dir / "artifacts").mkdir(exist_ok=True)
self._events_path = self.run_dir / "events.jsonl"
self._training_path = self.run_dir / "training_log.jsonl"
self._events_f = self._events_path.open("a", encoding="utf-8")
self._training_f = self._training_path.open("a", encoding="utf-8")
self._closed = False
def log_event(self, event: str, **payload: Any) -> None:
"""Write one event to events.jsonl and flush immediately."""
line = {"ts": _now_iso(), "event": event, **payload}
self._events_f.write(json.dumps(line) + "\n")
self._events_f.flush()
def log_training_step(self, lesson_idx: int, step: int, **fields: Any) -> None:
"""Write a per-step training record. Not flushed every line — this
channel is meant for high-frequency writes. Flushed on close()."""
line = {"ts": _now_iso(), "lesson_idx": lesson_idx, "step": step, **fields}
self._training_f.write(json.dumps(line) + "\n")
def close(self) -> None:
if self._closed:
return
self._training_f.flush()
self._training_f.close()
self._events_f.flush()
self._events_f.close()
self._closed = True
def __enter__(self) -> "RunRecorder":
return self
def __exit__(
self,
exc_type: type[BaseException] | None,
exc: BaseException | None,
tb: TracebackType | None,
) -> None:
self.close()