gabriel / musehub public
test_delta_compressed_base_regression.py python
510 lines 18.0 KB
Raw
sha256:9590cee1e0ccd6c76528f005b95d634d80f5019f0dcb7c371e149adc31d1fb65 refactor: enforce gRPC framing on all MWP wire traffic Sonnet 4.6 minor ⚠ breaking 156 days ago
1 """Regression tests: delta+zlib push when the base object is stored zlib-compressed.
2
3 Root cause:
4 Objects pushed via the old wire path were stored zlib-compressed in R2.
5 When a subsequent push sends a new object as delta+zlib against such a base,
6 wire_push_object_pack fetches the base via backend.get() — which returns the
7 raw R2 bytes (zlib-compressed), not the plain content. apply_delta() then
8 operates on compressed bytes and produces garbage. Hash verification is
9 skipped for non-"sha256:"-prefixed IDs, so the garbage passes silently.
10
11 Real-world failure:
12 staging.musehub.ai/gabriel/muse-zsh showed "0!" in the README after a one-line
13 README edit was pushed as delta+zlib against the old zlib-compressed README
14 object already in R2.
15
16 Fix:
17 1. musehub/muse_contracts/compression.py — expose decompress_if_needed().
18 2. musehub/services/musehub_wire.py — apply decompress_if_needed() to base_raw
19 before calling apply_delta() so the delta is always applied against plain bytes.
20 3. musehub/api/routes/wire.py + musehub/services/musehub_wire.py — add
21 POST /{owner}/{slug}/repair-object endpoint so operators can correct an object
22 already stored with wrong bytes without requiring direct DB/R2 access.
23 """
24 from __future__ import annotations
25
26 import hashlib
27 import struct
28 import zlib
29
30 import msgpack
31 import pytest
32 from httpx import AsyncClient
33 from sqlalchemy.ext.asyncio import AsyncSession
34
35 from tests.factories import create_repo as factory_create_repo
36 from musehub.db import musehub_models as db
37 from musehub.types.json_types import JSONValue, StrDict
38
39
40 # ---------------------------------------------------------------------------
41 # Helpers
42 # ---------------------------------------------------------------------------
43
44 def _sha256_oid(raw: bytes) -> str:
45 """Canonical sha256:-prefixed object_id — required by WireObject validator."""
46 return "sha256:" + hashlib.sha256(raw).hexdigest()
47
48
49 def _compress_zlib(raw: bytes) -> bytes:
50 return zlib.compress(raw)
51
52
53 def _make_delta_stream(target: bytes) -> bytes:
54 """Minimal DATA-only delta (no COPY), zlib-compressed.
55
56 Format: b'\\x01' + struct.pack('>I', len) + data
57 This is a valid encoding that always works regardless of the base content.
58 """
59 stream = b"\x01" + struct.pack(">I", len(target)) + target
60 return zlib.compress(stream, level=1)
61
62
63 def _mp(data: JSONValue) -> bytes:
64 return msgpack.packb(data, use_bin_type=True)
65
66
67 # ---------------------------------------------------------------------------
68 # Unit tests — decompress_if_needed in compression module
69 # ---------------------------------------------------------------------------
70
71 def test_decompress_if_needed_importable_from_compression_module() -> None:
72 """decompress_if_needed must live in muse_contracts.compression, not ui_tree."""
73 from musehub.types.compression import decompress_if_needed
74 assert callable(decompress_if_needed)
75
76
77 def test_decompress_if_needed_zlib_level_1_magic_78_01() -> None:
78 from musehub.types.compression import decompress_if_needed
79 data = zlib.compress(b"hello world\n", level=1)
80 assert data[:2] == b"\x78\x01"
81 assert decompress_if_needed(data) == b"hello world\n"
82
83
84 def test_decompress_if_needed_zlib_level_6_magic_78_9c() -> None:
85 from musehub.types.compression import decompress_if_needed
86 data = zlib.compress(b"hello world\n", level=6)
87 assert data[:2] == b"\x78\x9c"
88 assert decompress_if_needed(data) == b"hello world\n"
89
90
91 def test_decompress_if_needed_zlib_level_9_magic_78_da() -> None:
92 from musehub.types.compression import decompress_if_needed
93 data = zlib.compress(b"hello world\n", level=9)
94 assert data[:2] == b"\x78\xda"
95 assert decompress_if_needed(data) == b"hello world\n"
96
97
98 def test_decompress_if_needed_plain_text_passthrough() -> None:
99 from musehub.types.compression import decompress_if_needed
100 data = b"# plain README\n\nSome content.\n"
101 assert decompress_if_needed(data) == data
102
103
104 def test_decompress_if_needed_binary_non_zlib_passthrough() -> None:
105 from musehub.types.compression import decompress_if_needed
106 # PNG magic — must pass through unchanged
107 data = b"\x89PNG\r\n\x1a\n" + b"\x00" * 50
108 assert decompress_if_needed(data) == data
109
110
111 def test_decompress_if_needed_empty_bytes_passthrough() -> None:
112 from musehub.types.compression import decompress_if_needed
113 assert decompress_if_needed(b"") == b""
114
115
116 def test_decompress_if_needed_short_data_passthrough() -> None:
117 from musehub.types.compression import decompress_if_needed
118 assert decompress_if_needed(b"\x78") == b"\x78" # only 1 byte — no magic match
119
120
121 def test_decompress_if_needed_truncated_zlib_returns_original() -> None:
122 from musehub.types.compression import decompress_if_needed
123 # Valid header, invalid body — must not raise; returns original bytes.
124 data = b"\x78\x9c\x00"
125 result = decompress_if_needed(data)
126 assert isinstance(result, bytes)
127 assert result == data # original returned on decompression failure
128
129
130 def test_decompress_if_needed_full_readme_roundtrip() -> None:
131 from musehub.types.compression import decompress_if_needed
132 readme = b"# muse-zsh\n\nOh My ZSH plugin for Muse version control.\n" * 10
133 compressed = zlib.compress(readme)
134 assert decompress_if_needed(compressed) == readme
135
136
137 # ---------------------------------------------------------------------------
138 # Integration — baseline: delta against plain base works (before + after fix)
139 # ---------------------------------------------------------------------------
140
141 @pytest.mark.asyncio
142 async def test_push_delta_plain_base_baseline(
143 client: AsyncClient,
144 db_session: AsyncSession,
145 wire_headers: StrDict,
146 ) -> None:
147 """Baseline: delta+zlib against a plain-text base must work before and after fix."""
148 repo = await factory_create_repo(
149 db_session, slug="delta-plain-base-baseline", owner="test-user-wire"
150 )
151 path = "README.md"
152 base_raw = b"# Repo\n\nInitial content.\n" * 30
153 base_oid = _sha256_oid(base_raw)
154
155 r_base = await client.post(
156 f"/{repo.owner}/{repo.slug}/push/object-pack",
157 content=_mp({"objects": [
158 {"object_id": base_oid, "content": base_raw, "path": path, "encoding": "raw"}
159 ]}),
160 headers=wire_headers,
161 )
162 assert r_base.status_code == 200
163
164 target_raw = b"# Repo\n\nUpdated content.\n" * 30
165 target_oid = _sha256_oid(target_raw)
166
167 r_target = await client.post(
168 f"/{repo.owner}/{repo.slug}/push/object-pack",
169 content=_mp({"objects": [{
170 "object_id": target_oid,
171 "content": _make_delta_stream(target_raw),
172 "path": path,
173 "encoding": "delta+zlib",
174 "base_id": base_oid,
175 }]}),
176 headers=wire_headers,
177 )
178 assert r_target.status_code == 200
179 assert msgpack.unpackb(r_target.content, raw=False)["stored"] == 1
180
181
182 # ---------------------------------------------------------------------------
183 # Integration — regression: delta against ZLIB-COMPRESSED base
184 # ---------------------------------------------------------------------------
185
186 @pytest.mark.asyncio
187 async def test_push_delta_zlib_base_in_content_cache_succeeds(
188 client: AsyncClient,
189 db_session: AsyncSession,
190 wire_headers: StrDict,
191 ) -> None:
192 """delta+zlib push must succeed when the base is in content_cache as zlib bytes.
193
194 Simulates the old wire path: base object arrived zlib-compressed and was stored
195 in content_cache without decompression. The server must decompress before
196 apply_delta(); otherwise the delta is applied against the compressed bytes and
197 the result is garbage.
198 """
199 repo = await factory_create_repo(
200 db_session, slug="delta-zlib-cache-base", owner="test-user-wire"
201 )
202 path = "README.md"
203
204 base_raw = b"# muse-zsh\n\nOh My ZSH plugin.\n" * 30
205 base_oid = _sha256_oid(base_raw)
206 base_zlib = _compress_zlib(base_raw)
207
208 # Inject base row with ZLIB bytes in content_cache — simulating old wire path.
209 await db_session.execute(
210 db.MusehubObject.__table__.insert().values(
211 object_id=base_oid,
212 path=path,
213 size_bytes=len(base_raw),
214 disk_path="",
215 storage_uri="pending",
216 content_cache=base_zlib, # ← zlib bytes, not plain
217 )
218 )
219 db_session.add(db.MusehubObjectRef(repo_id=repo.repo_id, object_id=base_oid))
220 await db_session.commit()
221
222 target_raw = b"# muse-zsh\n\nUpdated plugin.\n" * 30
223 target_oid = _sha256_oid(target_raw)
224
225 r = await client.post(
226 f"/{repo.owner}/{repo.slug}/push/object-pack",
227 content=_mp({"objects": [{
228 "object_id": target_oid,
229 "content": _make_delta_stream(target_raw),
230 "path": path,
231 "encoding": "delta+zlib",
232 "base_id": base_oid,
233 }]}),
234 headers=wire_headers,
235 )
236 assert r.status_code == 200, r.text
237 assert msgpack.unpackb(r.content, raw=False)["stored"] == 1
238
239
240 @pytest.mark.asyncio
241 async def test_push_delta_zlib_base_in_storage_succeeds(
242 client: AsyncClient,
243 db_session: AsyncSession,
244 wire_headers: StrDict,
245 ) -> None:
246 """delta+zlib push must succeed when the base is in storage as zlib bytes.
247
248 Simulates R2: base object was uploaded as zlib-compressed (old wire path),
249 content_cache is NULL. backend.get() returns zlib bytes; server must
250 decompress before apply_delta().
251 """
252 import musehub.services.musehub_wire as _wire_svc
253
254 repo = await factory_create_repo(
255 db_session, slug="delta-zlib-storage-base", owner="test-user-wire"
256 )
257 path = "README.md"
258
259 base_raw = b"# muse-zsh\n\nOld content.\n" * 30
260 base_oid = _sha256_oid(base_raw)
261 base_zlib = _compress_zlib(base_raw)
262
263 backend = _wire_svc.get_backend()
264 # Store ZLIB-COMPRESSED bytes in storage — simulating what old push path did.
265 await backend.put(base_oid, base_zlib)
266
267 await db_session.execute(
268 db.MusehubObject.__table__.insert().values(
269 object_id=base_oid,
270 path=path,
271 size_bytes=len(base_raw),
272 disk_path="",
273 storage_uri=backend.uri_for(base_oid),
274 content_cache=None, # ← forces fetch from storage backend
275 )
276 )
277 db_session.add(db.MusehubObjectRef(repo_id=repo.repo_id, object_id=base_oid))
278 await db_session.commit()
279
280 target_raw = b"# muse-zsh\n\nNew content.\n" * 30
281 target_oid = _sha256_oid(target_raw)
282
283 r = await client.post(
284 f"/{repo.owner}/{repo.slug}/push/object-pack",
285 content=_mp({"objects": [{
286 "object_id": target_oid,
287 "content": _make_delta_stream(target_raw),
288 "path": path,
289 "encoding": "delta+zlib",
290 "base_id": base_oid,
291 }]}),
292 headers=wire_headers,
293 )
294 assert r.status_code == 200, r.text
295 assert msgpack.unpackb(r.content, raw=False)["stored"] == 1
296
297
298 @pytest.mark.asyncio
299 async def test_push_delta_reconstructed_bytes_match_target(
300 client: AsyncClient,
301 db_session: AsyncSession,
302 wire_headers: StrDict,
303 ) -> None:
304 """The bytes written to storage after delta reconstruction must equal the original target.
305
306 The critical assertion: stored != garbage, stored == target_raw.
307 """
308 import musehub.services.musehub_wire as _wire_svc
309
310 repo = await factory_create_repo(
311 db_session, slug="delta-content-exact", owner="test-user-wire"
312 )
313 path = "README.md"
314
315 base_raw = b"# muse-zsh\n\nOld.\n" * 50
316 base_oid = _sha256_oid(base_raw)
317 base_zlib = _compress_zlib(base_raw)
318
319 backend = _wire_svc.get_backend()
320 await backend.put(base_oid, base_zlib)
321 await db_session.execute(
322 db.MusehubObject.__table__.insert().values(
323 object_id=base_oid, path=path,
324 size_bytes=len(base_raw), disk_path="",
325 storage_uri=backend.uri_for(base_oid),
326 content_cache=None,
327 )
328 )
329 db_session.add(db.MusehubObjectRef(repo_id=repo.repo_id, object_id=base_oid))
330 await db_session.commit()
331
332 target_raw = b"# muse-zsh\n\nNew.\n" * 50
333 target_oid = _sha256_oid(target_raw)
334
335 r = await client.post(
336 f"/{repo.owner}/{repo.slug}/push/object-pack",
337 content=_mp({"objects": [{
338 "object_id": target_oid,
339 "content": _make_delta_stream(target_raw),
340 "path": path,
341 "encoding": "delta+zlib",
342 "base_id": base_oid,
343 }]}),
344 headers=wire_headers,
345 )
346 assert r.status_code == 200
347
348 stored = await backend.get(target_oid)
349 assert stored is not None, "target object must be written to storage"
350 assert stored == target_raw, (
351 f"stored bytes do not match target:\n"
352 f" first 80 stored: {stored[:80]!r}\n"
353 f" first 80 target: {target_raw[:80]!r}"
354 )
355
356
357 # ---------------------------------------------------------------------------
358 # repair-object endpoint
359 # ---------------------------------------------------------------------------
360
361 @pytest.mark.asyncio
362 async def test_repair_object_corrects_corrupted_stored_bytes(
363 client: AsyncClient,
364 db_session: AsyncSession,
365 wire_headers: StrDict,
366 ) -> None:
367 """POST /repair-object replaces corrupt stored bytes with the verified correct content.
368
369 Simulates the staging scenario: an object has garbage bytes in storage (from a
370 failed delta reconstruction). repair-object accepts the correct raw bytes,
371 verifies SHA-256, overwrites storage, and updates the DB row.
372 """
373 import musehub.services.musehub_wire as _wire_svc
374
375 repo = await factory_create_repo(
376 db_session, slug="repair-corrupted-object", owner="test-user-wire"
377 )
378 correct_raw = b"# README\n\nCorrect content.\n" * 20
379 object_id = _sha256_oid(correct_raw)
380 garbage = b"\x00\x01\x02\x03garbage\xff"
381
382 backend = _wire_svc.get_backend()
383 await backend.put(object_id, garbage)
384 await db_session.execute(
385 db.MusehubObject.__table__.insert().values(
386 object_id=object_id,
387 path="README.md",
388 size_bytes=len(garbage),
389 disk_path="",
390 storage_uri=backend.uri_for(object_id),
391 content_cache=None,
392 )
393 )
394 db_session.add(db.MusehubObjectRef(repo_id=repo.repo_id, object_id=object_id))
395 await db_session.commit()
396
397 r = await client.post(
398 f"/{repo.owner}/{repo.slug}/repair-object",
399 content=_mp({"object_id": object_id, "content": correct_raw}),
400 headers=wire_headers,
401 )
402 assert r.status_code == 200, r.text
403 data = msgpack.unpackb(r.content, raw=False)
404 assert data["repaired"] is True
405
406 stored = await backend.get(object_id)
407 assert stored == correct_raw
408
409
410 @pytest.mark.asyncio
411 async def test_repair_object_rejects_hash_mismatch(
412 client: AsyncClient,
413 db_session: AsyncSession,
414 wire_headers: StrDict,
415 ) -> None:
416 """repair-object rejects content whose SHA-256 does not match the declared object_id."""
417 repo = await factory_create_repo(
418 db_session, slug="repair-hash-mismatch", owner="test-user-wire"
419 )
420 correct_raw = b"correct content bytes for hash test"
421 wrong_content = b"completely different content"
422 object_id = _sha256_oid(correct_raw)
423
424 r = await client.post(
425 f"/{repo.owner}/{repo.slug}/repair-object",
426 content=_mp({"object_id": object_id, "content": wrong_content}),
427 headers=wire_headers,
428 )
429 assert r.status_code == 422, r.text
430
431
432 @pytest.mark.asyncio
433 async def test_repair_object_accepts_sha256_prefixed_id(
434 client: AsyncClient,
435 db_session: AsyncSession,
436 wire_headers: StrDict,
437 ) -> None:
438 """repair-object works with sha256:-prefixed IDs."""
439 import musehub.services.musehub_wire as _wire_svc
440
441 repo = await factory_create_repo(
442 db_session, slug="repair-prefixed-id", owner="test-user-wire"
443 )
444 correct_raw = b"content with prefixed object_id\n" * 10
445 object_id_prefixed = _sha256_oid(correct_raw) # "sha256:<hex>"
446
447 backend = _wire_svc.get_backend()
448 await backend.put(object_id_prefixed, b"garbage")
449 await db_session.execute(
450 db.MusehubObject.__table__.insert().values(
451 object_id=object_id_prefixed,
452 path="src/file.py",
453 size_bytes=7,
454 disk_path="",
455 storage_uri=backend.uri_for(object_id_prefixed),
456 content_cache=None,
457 )
458 )
459 db_session.add(db.MusehubObjectRef(repo_id=repo.repo_id, object_id=object_id_prefixed))
460 await db_session.commit()
461
462 r = await client.post(
463 f"/{repo.owner}/{repo.slug}/repair-object",
464 content=_mp({"object_id": object_id_prefixed, "content": correct_raw}),
465 headers=wire_headers,
466 )
467 assert r.status_code == 200, r.text
468 assert msgpack.unpackb(r.content, raw=False)["repaired"] is True
469
470
471 @pytest.mark.asyncio
472 async def test_repair_object_is_idempotent(
473 client: AsyncClient,
474 db_session: AsyncSession,
475 wire_headers: StrDict,
476 ) -> None:
477 """Calling repair-object twice with the same correct content is safe."""
478 import musehub.services.musehub_wire as _wire_svc
479
480 repo = await factory_create_repo(
481 db_session, slug="repair-idempotent-test", owner="test-user-wire"
482 )
483 correct_raw = b"idempotent repair content\n" * 15
484 object_id = _sha256_oid(correct_raw)
485
486 backend = _wire_svc.get_backend()
487 await backend.put(object_id, b"garbage bytes")
488 await db_session.execute(
489 db.MusehubObject.__table__.insert().values(
490 object_id=object_id,
491 path="README.md",
492 size_bytes=13,
493 disk_path="",
494 storage_uri=backend.uri_for(object_id),
495 content_cache=None,
496 )
497 )
498 db_session.add(db.MusehubObjectRef(repo_id=repo.repo_id, object_id=object_id))
499 await db_session.commit()
500
501 url = f"/{repo.owner}/{repo.slug}/repair-object"
502 payload = _mp({"object_id": object_id, "content": correct_raw})
503
504 r1 = await client.post(url, content=payload, headers=wire_headers)
505 assert r1.status_code == 200
506 r2 = await client.post(url, content=payload, headers=wire_headers)
507 assert r2.status_code == 200
508
509 stored = await backend.get(object_id)
510 assert stored == correct_raw
File History 1 commit
sha256:9590cee1e0ccd6c76528f005b95d634d80f5019f0dcb7c371e149adc31d1fb65 refactor: enforce gRPC framing on all MWP wire traffic Sonnet 4.6 minor ⚠ 156 days ago