test_wire_snapshot_delta.py
python
sha256:b636f72dcba9e190afb980bece906fa5b717fbde014b76ef023df8cb96e01eb9
docs: expand cache plan with all seven testing tiers and do…
Sonnet 4.6
131 days ago
| 1 | """Tests for snapshot delta encoding in _frame_generator. |
| 2 | |
| 3 | Verifies that _frame_generator: |
| 4 | A. Sends the first snapshot as a full manifest |
| 5 | B. Sends subsequent snapshots as deltas |
| 6 | C. Deltas chain correctly (each references the previous snapshot) |
| 7 | D. A stream with only one commit sends a full snapshot (no delta) |
| 8 | E. write_commit_pack includes snapshot_deltas in msgpack output |
| 9 | F. Applying every delta in stream order reconstructs all manifests |
| 10 | """ |
| 11 | from __future__ import annotations |
| 12 | |
| 13 | import struct |
| 14 | import msgpack |
| 15 | from muse.core._types import Manifest, MsgpackDict, blob_id |
| 16 | |
| 17 | |
| 18 | # --------------------------------------------------------------------------- |
| 19 | # Helpers |
| 20 | # --------------------------------------------------------------------------- |
| 21 | |
| 22 | def _parse_mwp_frames(data: bytes) -> list[dict]: |
| 23 | """Parse concatenated MWP-framed bytes into decoded payload dicts.""" |
| 24 | frames: list[dict] = [] |
| 25 | pos = 0 |
| 26 | magic = b"muse" |
| 27 | while pos < len(data): |
| 28 | assert data[pos:pos + 4] == magic, f"bad magic at {pos}" |
| 29 | pos += 4 |
| 30 | pos += 1 # version byte |
| 31 | header_len = struct.unpack_from(">I", data, pos)[0] |
| 32 | pos += 4 + header_len |
| 33 | payload_len = struct.unpack_from(">Q", data, pos)[0] |
| 34 | pos += 8 |
| 35 | payload_bytes = data[pos:pos + payload_len] |
| 36 | pos += payload_len |
| 37 | frames.append(msgpack.unpackb(payload_bytes, raw=False)) |
| 38 | return frames |
| 39 | |
| 40 | |
| 41 | def _parse_mwp_payload(mwp_frame: bytes) -> MsgpackDict: |
| 42 | """Extract msgpack payload from a raw MWP frame (past the binary envelope).""" |
| 43 | # MWP: b"muse" | 0x01 | uint32 header_len | msgpack(envelope) | uint64 payload_len | payload |
| 44 | assert mwp_frame[:4] == b"muse" |
| 45 | assert mwp_frame[4] == 0x01 |
| 46 | header_len = struct.unpack(">I", mwp_frame[5:9])[0] |
| 47 | payload_offset = 9 + header_len + 8 # skip envelope + 8-byte payload_len |
| 48 | payload = mwp_frame[payload_offset:] |
| 49 | return msgpack.unpackb(payload, raw=False) |
| 50 | |
| 51 | |
| 52 | def _build_manifest(n_files: int) -> Manifest: |
| 53 | return {f"src/module_{i}.py": f"sha256:{'a' * 60}{i:04d}" for i in range(n_files)} |
| 54 | |
| 55 | |
| 56 | def _make_commit(snapshot_id: str, i: int) -> MsgpackDict: |
| 57 | return { |
| 58 | "commit_id": f"sha256:{'c' * 60}{i:04d}", |
| 59 | "snapshot_id": snapshot_id, |
| 60 | "message": f"commit {i}", |
| 61 | "branch": "dev", |
| 62 | "author": "gabriel", |
| 63 | "committed_at": "2026-04-23T00:00:00+00:00", |
| 64 | "parent_commit_id": None, |
| 65 | "agent_id": "", |
| 66 | "model_id": "", |
| 67 | "toolchain_id": "", |
| 68 | "signer_key": "", |
| 69 | "signature": "", |
| 70 | } |
| 71 | |
| 72 | |
| 73 | def _make_snapshot(manifest: Manifest, i: int) -> MsgpackDict: |
| 74 | snap_bytes = msgpack.packb(manifest, use_bin_type=True) |
| 75 | snap_id = blob_id(snap_bytes) |
| 76 | return { |
| 77 | "snapshot_id": snap_id, |
| 78 | "manifest": manifest, |
| 79 | "directories": [], |
| 80 | "created_at": "2026-04-23T00:00:00+00:00", |
| 81 | "note": "", |
| 82 | } |
| 83 | |
| 84 | |
| 85 | def _build_linear_history(n_commits: int, n_files: int) -> tuple[list, list, dict]: |
| 86 | """Build a sequence of commits each changing 3 files in a manifest. |
| 87 | |
| 88 | Returns: (commits, snapshots, snap_by_id) |
| 89 | """ |
| 90 | commits = [] |
| 91 | snapshots = [] |
| 92 | snap_by_id = {} |
| 93 | |
| 94 | manifest = _build_manifest(n_files) |
| 95 | for i in range(n_commits): |
| 96 | # Mutate 3 files per commit |
| 97 | new_manifest = dict(manifest) |
| 98 | for j in range(3): |
| 99 | path = f"src/module_{(i * 3 + j) % n_files}.py" |
| 100 | new_manifest[path] = f"sha256:{'b' * 60}{i * 3 + j:04d}" |
| 101 | snap = _make_snapshot(new_manifest, i) |
| 102 | commit = _make_commit(snap["snapshot_id"], i) |
| 103 | commits.append(commit) |
| 104 | snapshots.append(snap) |
| 105 | snap_by_id[snap["snapshot_id"]] = snap |
| 106 | manifest = new_manifest |
| 107 | |
| 108 | return commits, snapshots, snap_by_id |
| 109 | |
| 110 | |
| 111 | # --------------------------------------------------------------------------- |
| 112 | # A. First snapshot in stream is always a full manifest |
| 113 | # --------------------------------------------------------------------------- |
| 114 | |
| 115 | def test_first_snapshot_is_full(): |
| 116 | from muse.core.transport import _frame_generator |
| 117 | |
| 118 | commits, snapshots, snap_by_id = _build_linear_history(n_commits=5, n_files=100) |
| 119 | |
| 120 | raw = b"".join(_frame_generator([], commits, snapshots, local_head=None)) |
| 121 | c_frames = [f for f in _parse_mwp_frames(raw) if f.get("t") == "C"] |
| 122 | assert c_frames, "expected at least one C frame" |
| 123 | assert len(c_frames[0].get("snapshots", [])) >= 1, "first C frame must contain at least one full snapshot" |
| 124 | |
| 125 | |
| 126 | # --------------------------------------------------------------------------- |
| 127 | # B. Subsequent snapshots sent as deltas |
| 128 | # --------------------------------------------------------------------------- |
| 129 | |
| 130 | def test_subsequent_snapshots_are_deltas(): |
| 131 | from muse.core.transport import _frame_generator |
| 132 | |
| 133 | commits, snapshots, _ = _build_linear_history(n_commits=10, n_files=100) |
| 134 | |
| 135 | raw = b"".join(_frame_generator([], commits, snapshots, local_head=None)) |
| 136 | c_frames = [f for f in _parse_mwp_frames(raw) if f.get("t") == "C"] |
| 137 | |
| 138 | all_deltas = [] |
| 139 | for cf in c_frames: |
| 140 | all_deltas.extend(cf.get("snapshot_deltas") or []) |
| 141 | |
| 142 | assert len(all_deltas) > 0, "expected delta snapshots for commits 2+" |
| 143 | |
| 144 | |
| 145 | # --------------------------------------------------------------------------- |
| 146 | # C. Delta chain integrity — each delta references a previously sent snapshot |
| 147 | # --------------------------------------------------------------------------- |
| 148 | |
| 149 | def test_delta_chain_references_valid_base(): |
| 150 | from muse.core.transport import _frame_generator |
| 151 | |
| 152 | commits, snapshots, _ = _build_linear_history(n_commits=20, n_files=100) |
| 153 | |
| 154 | raw = b"".join(_frame_generator([], commits, snapshots, local_head=None)) |
| 155 | c_frames = [f for f in _parse_mwp_frames(raw) if f.get("t") == "C"] |
| 156 | |
| 157 | seen_snapshot_ids: set[str] = set() |
| 158 | for cf in c_frames: |
| 159 | for snap in cf.get("snapshots") or []: |
| 160 | seen_snapshot_ids.add(snap["snapshot_id"]) |
| 161 | for delta in cf.get("snapshot_deltas") or []: |
| 162 | base_id = delta.get("base_id") |
| 163 | assert base_id in seen_snapshot_ids, ( |
| 164 | f"delta references base {base_id!r} not yet seen in stream" |
| 165 | ) |
| 166 | seen_snapshot_ids.add(delta["snapshot_id"]) |
| 167 | |
| 168 | |
| 169 | # --------------------------------------------------------------------------- |
| 170 | # D. Single commit — full snapshot, no deltas |
| 171 | # --------------------------------------------------------------------------- |
| 172 | |
| 173 | def test_single_commit_sends_full_snapshot_no_deltas(): |
| 174 | from muse.core.transport import _frame_generator |
| 175 | |
| 176 | commits, snapshots, _ = _build_linear_history(n_commits=1, n_files=50) |
| 177 | |
| 178 | raw = b"".join(_frame_generator([], commits, snapshots, local_head=None)) |
| 179 | c_frames = [f for f in _parse_mwp_frames(raw) if f.get("t") == "C"] |
| 180 | |
| 181 | assert len(c_frames) == 1 |
| 182 | assert len(c_frames[0].get("snapshots", [])) == 1 |
| 183 | assert c_frames[0].get("snapshot_deltas") in (None, []) |
| 184 | |
| 185 | |
| 186 | # --------------------------------------------------------------------------- |
| 187 | # E. write_commit_pack includes snapshot_deltas key when provided |
| 188 | # --------------------------------------------------------------------------- |
| 189 | |
| 190 | def test_write_commit_pack_includes_snapshot_deltas(): |
| 191 | from muse.core.mpack import MPackStreamWriter |
| 192 | writer = MPackStreamWriter() |
| 193 | |
| 194 | delta = { |
| 195 | "snapshot_id": "sha256:" + "a" * 64, |
| 196 | "base_id": "sha256:" + "b" * 64, |
| 197 | "added": {"src/foo.py": "sha256:" + "c" * 64}, |
| 198 | "removed": [], |
| 199 | "directories": [], |
| 200 | "created_at": "2026-04-23T00:00:00+00:00", |
| 201 | "note": "", |
| 202 | } |
| 203 | raw = writer.write_commit_pack(commits=[], snapshots=[], snapshot_deltas=[delta]) |
| 204 | frame = msgpack.unpackb(raw, raw=False) |
| 205 | assert "snapshot_deltas" in frame |
| 206 | assert len(frame["snapshot_deltas"]) == 1 |
| 207 | assert frame["snapshot_deltas"][0]["base_id"] == "sha256:" + "b" * 64 |
| 208 | |
| 209 | |
| 210 | def test_write_commit_pack_no_deltas_key_when_empty(): |
| 211 | from muse.core.mpack import MPackStreamWriter |
| 212 | writer = MPackStreamWriter() |
| 213 | raw = writer.write_commit_pack(commits=[], snapshots=[]) |
| 214 | frame = msgpack.unpackb(raw, raw=False) |
| 215 | assert "snapshot_deltas" not in frame |
| 216 | |
| 217 | |
| 218 | # --------------------------------------------------------------------------- |
| 219 | # F. Applying every delta in stream order reconstructs all manifests |
| 220 | # --------------------------------------------------------------------------- |
| 221 | |
| 222 | def test_delta_chain_full_reconstruction(): |
| 223 | """All manifests reconstructed by walking the delta chain must equal |
| 224 | the original snapshots passed to _frame_generator.""" |
| 225 | from muse.core.transport import _frame_generator |
| 226 | from muse.core.mpack import apply_snapshot_delta |
| 227 | |
| 228 | n_commits = 30 |
| 229 | commits, snapshots, snap_by_id = _build_linear_history(n_commits=n_commits, n_files=200) |
| 230 | |
| 231 | raw = b"".join(_frame_generator([], commits, snapshots, local_head=None)) |
| 232 | c_frames = [f for f in _parse_mwp_frames(raw) if f.get("t") == "C"] |
| 233 | |
| 234 | # Walk stream, reconstructing every snapshot |
| 235 | reconstructed: dict[str, dict] = {} # snapshot_id → manifest |
| 236 | |
| 237 | for cf in c_frames: |
| 238 | for snap in cf.get("snapshots") or []: |
| 239 | reconstructed[snap["snapshot_id"]] = snap["manifest"] |
| 240 | for delta in cf.get("snapshot_deltas") or []: |
| 241 | base_manifest = reconstructed[delta["base_id"]] |
| 242 | full = apply_snapshot_delta(base_manifest, delta["added"], delta["removed"]) |
| 243 | reconstructed[delta["snapshot_id"]] = full |
| 244 | |
| 245 | for snap in snapshots: |
| 246 | sid = snap["snapshot_id"] |
| 247 | assert sid in reconstructed, f"snapshot {sid[:16]} not in reconstructed" |
| 248 | assert reconstructed[sid] == snap["manifest"], ( |
| 249 | f"reconstructed manifest for {sid[:16]} does not match original" |
| 250 | ) |
File History
1 commit
sha256:b636f72dcba9e190afb980bece906fa5b717fbde014b76ef023df8cb96e01eb9
docs: expand cache plan with all seven testing tiers and do…
Sonnet 4.6
131 days ago