"""A decision log reads the same compressed as it does plain. A finished log compresses about 30 to 1 - `phi` writes all forty feature names again for every candidate, and a 300-match corpus holds millions of them - so 53GB of corpus becomes 3.2GB. That is the difference between keeping the rows a refit needs and deleting them for space, which is a choice this project has already made the wrong way once. `sds.match` writes the compressed form, as each match's container exits, and everything reads it. That is the half these tests grew: a recording run that only stayed inside the disk because somebody watched it and ran `gzip` by hand is not a recording run anybody can leave alone, and one reached 17G and 98% of this machine's disk while it was still playing. """ from __future__ import annotations import gzip import json import tempfile import unittest from pathlib import Path from sds.match import compress_decision_logs, compress_log from sds.viewer import COMPRESSED_SUFFIX, DECISION_SUFFIX, _tag_of, decision_logs, iter_decisions ROWS = [ {"seq": 0, "round": 1, "phase": "movement", "candidates": [{"phi": {"cohesion": 1.0}}]}, {"seq": 1, "round": 1, "phase": "firing", "candidates": [{"phi": {"cohesion": 0.5}}]}, ] class TestCompressedLogs(unittest.TestCase): def setUp(self): self.dir = Path(tempfile.mkdtemp()) body = "".join(json.dumps(r) + "\n" for r in ROWS) self.plain = self.dir / f"tag-seed1-abcd1234-sds_North{DECISION_SUFFIX}" self.plain.write_text(body) self.gz = self.dir / f"tag-seed2-abcd5678-sds_North{COMPRESSED_SUFFIX}" with gzip.open(self.gz, "wt", encoding="utf-8") as handle: handle.write(body) def test_the_two_forms_read_the_same_records(self): self.assertEqual(list(iter_decisions(self.plain)), list(iter_decisions(self.gz))) def test_both_forms_are_discovered(self): self.assertEqual(decision_logs(self.dir), sorted([self.gz, self.plain])) def test_a_corpus_may_hold_a_mix(self): # Which is the normal state part-way through a compression pass. found = decision_logs(self.dir) self.assertEqual(len(found), 2) rows = [r for log in found for r in iter_decisions(log)] self.assertEqual(len(rows), 2 * len(ROWS)) def test_the_tag_survives_the_extra_suffix(self): # `.gz` is stripped before the seat is split off, or the tag keeps it # and a compressed log stops matching its own result document. self.assertEqual(_tag_of(self.gz), "tag-seed2-abcd5678") self.assertEqual(_tag_of(self.plain), "tag-seed1-abcd1234") def test_a_truncated_member_yields_what_it_had(self): # The plain reader has always tolerated a half-written last line. A # compressed one raises instead of returning a short read, so that is # caught too rather than losing the whole log. cut = self.dir / f"cut-seed3-abcd9999-sds_North{COMPRESSED_SUFFIX}" cut.write_bytes(self.gz.read_bytes()[:-6]) self.assertLessEqual(len(list(iter_decisions(cut))), len(ROWS)) class TestCompressingWhatAMatchWrote(unittest.TestCase): """`sds.match` compressing one match's logs when its container exits.""" def setUp(self): self.dir = Path(tempfile.mkdtemp()) self.body = "".join(json.dumps(r) + "\n" for r in ROWS) * 200 self.tag = "tag-seed1-abcd1234" self.logs = [] for seat in ("sds_North", "sds_South"): log = self.dir / f"{self.tag}-{seat}{DECISION_SUFFIX}" log.write_text(self.body) self.logs.append(log) def test_the_bytes_come_back_exactly(self): # The whole safety argument. If the decompressed file is the file that # was there, no instrument can tell the difference and no field can # have been dropped. packed = compress_log(self.logs[0]) self.assertEqual(gzip.decompress(packed.read_bytes()).decode(), self.body) self.assertFalse(self.logs[0].exists()) def test_it_gets_smaller(self): before = self.logs[0].stat().st_size self.assertLess(compress_log(self.logs[0]).stat().st_size, before) def test_the_same_log_compresses_to_the_same_bytes(self): # `mtime=0` and no stored filename. Two runs of the same match must not # differ by a header nothing reads. other = self.dir / f"copy-seed1-abcd1234-sds_North{DECISION_SUFFIX}" other.write_text(self.body) self.assertEqual(compress_log(self.logs[0]).read_bytes(), compress_log(other).read_bytes()) def test_only_this_match_is_touched(self): # A benchmark plays several matches at once into one run directory, and # an agent in another worktree has its own. Neither may be reached. neighbour = self.dir / f"other-seed9-99999999-sds_North{DECISION_SUFFIX}" neighbour.write_text(self.body) compress_decision_logs(self.dir, self.tag) self.assertTrue(neighbour.exists()) self.assertEqual(len(list(self.dir.glob("*" + COMPRESSED_SUFFIX))), 2) def test_it_reports_what_it_saved(self): before = sum(log.stat().st_size for log in self.logs) saved = compress_decision_logs(self.dir, self.tag) after = sum(p.stat().st_size for p in self.dir.glob("*" + COMPRESSED_SUFFIX)) self.assertEqual(saved, before - after) def test_nothing_partial_is_left_behind(self): compress_decision_logs(self.dir, self.tag) self.assertEqual(list(self.dir.glob("*.part")), []) def test_the_records_survive_the_round_trip(self): wanted = list(iter_decisions(self.logs[0])) compress_decision_logs(self.dir, self.tag) packed = self.dir / f"{self.tag}-sds_North{COMPRESSED_SUFFIX}" self.assertEqual(list(iter_decisions(packed)), wanted) class TestEveryReaderOpensBothForms(unittest.TestCase): """The readers that globbed the plain name only, and said nothing. Compressing a run as it plays is only safe if everything that reads a run reads the compressed form. Three of these did not, and two of the three failed by reporting *zero* rather than by raising - `sds spread` opened a `.gz` as text, matched no line, and printed a table over "0 decisions". """ def setUp(self): self.dir = Path(tempfile.mkdtemp()) rows = [ { "seq": 0, "round": 1, "phase": "MOVEMENT", "candidates": [ {"label": "a", "phi": {"cohesion": 1.0}, "value": 1.0}, {"label": "b", "phi": {"cohesion": 0.0}, "value": 0.0}, ], "chosen": 0, }, { "seq": 1, "round": 1, "phase": "FIRING", "candidates": [{"label": "c", "phi": {"cohesion": 0.5}, "value": 0.5}], "chosen": 0, "firing": {"offered": [], "declared": []}, }, ] body = "".join(json.dumps(r) + "\n" for r in rows) self.tag = "tag-seed1-abcd1234" self.plain = self.dir / f"{self.tag}-sds_North{DECISION_SUFFIX}" self.plain.write_text(body) (self.dir / f"{self.tag}-sds_North.weights.json").write_text( json.dumps({"weights": {"cohesion": 1.0}}) ) self.packed = compress_log(self.dir / f"{self.tag}-sds_North{DECISION_SUFFIX}") self.plain = self.dir / f"{self.tag}-sds_North{DECISION_SUFFIX}" self.plain.write_text(body) def test_spread_counts_the_same_decisions(self): from sds import spread both = spread.scan(self.dir) self.assertEqual(both.logs, 2) self.assertEqual(both.rows, 4) def test_explain_finds_and_reads_a_compressed_log(self): from sds import explain self.assertIn(self.packed, explain.decision_logs(self.dir)) self.assertEqual(len(explain.read(self.packed)), 2) def test_explain_finds_the_sidecar_beside_a_compressed_log(self): # `name.replace(DECISION_SUFFIX, ...)` on a `.gz` asks for # `*.weights.json.gz`, which is not a file anybody writes. from sds import explain self.assertEqual(explain.weights_for(self.packed), {"cohesion": 1.0}) def test_firingaudit_finds_the_firing_rows(self): from sds import firingaudit self.assertEqual(len(firingaudit.rows(self.dir)), 2) def test_watch_replays_a_compressed_log(self): # `sds watch` on a finished run is a replay, and a `.gz` read as raw # bytes decodes to nothing at all rather than to a short read. from sds.live import LogTail tail = LogTail(self.packed, "sds_North") first = tail.poll() self.assertEqual(len(first), 2) self.assertEqual(tail.poll(), []) self.assertEqual(first[0]["seat"], "sds_North") if __name__ == "__main__": unittest.main()