test_object_store.py
python
sha256:a34090cc4a394a78bd72cbbe34b08cc59525141e19135b6c0ab154f10611b9ef
debug(push/stream): instrument O-frame decode path with INF…
Sonnet 4.6
patch
121 days ago
| 1 | """Tests for checklist section 5.2 — Object store. |
| 2 | |
| 3 | Covers: |
| 4 | - Object files are immutable after write (LocalBackend chmod 0o444) |
| 5 | - Integrity scan detects missing objects |
| 6 | - Integrity scan detects hash-mismatched objects |
| 7 | - Integrity scan passes on clean objects |
| 8 | - Per-repo storage quota enforced in wire push |
| 9 | - Soft-delete sets deleted_at without removing the file |
| 10 | - Hard-delete reaper removes only objects past retention window |
| 11 | """ |
| 12 | from __future__ import annotations |
| 13 | |
| 14 | import stat |
| 15 | import tempfile |
| 16 | from pathlib import Path |
| 17 | from unittest.mock import AsyncMock, MagicMock, patch |
| 18 | |
| 19 | import pytest |
| 20 | from sqlalchemy.ext.asyncio import AsyncSession |
| 21 | |
| 22 | from muse.core.types import blob_id, fake_id |
| 23 | from musehub.db.musehub_models import MusehubObject as _MusehubObject |
| 24 | from musehub.types.json_types import JSONValue |
| 25 | from tests.factories import create_repo |
| 26 | |
| 27 | |
| 28 | # --------------------------------------------------------------------------- |
| 29 | # Helpers |
| 30 | # --------------------------------------------------------------------------- |
| 31 | |
| 32 | async def _add_object( |
| 33 | session: AsyncSession, |
| 34 | repo_id: str, |
| 35 | data: bytes, |
| 36 | *, |
| 37 | object_id: str | None = None, |
| 38 | ) -> _MusehubObject: |
| 39 | from musehub.db import musehub_models as db |
| 40 | |
| 41 | oid = object_id or blob_id(data) |
| 42 | obj = db.MusehubObject( |
| 43 | object_id=oid, |
| 44 | path="test.bin", |
| 45 | size_bytes=len(data), |
| 46 | disk_path=f"/tmp/{oid.replace(':', '_')}", |
| 47 | storage_uri=f"local:///tmp/{oid.replace(':', '_')}", |
| 48 | ) |
| 49 | session.add(obj) |
| 50 | session.add(db.MusehubObjectRef(repo_id=repo_id, object_id=oid)) |
| 51 | await session.flush() |
| 52 | return obj |
| 53 | |
| 54 | |
| 55 | # --------------------------------------------------------------------------- |
| 56 | # Immutability — LocalBackend sets file to read-only after write |
| 57 | # --------------------------------------------------------------------------- |
| 58 | |
| 59 | def test_local_backend_write_sets_readonly() -> None: |
| 60 | """LocalBackend._write must chmod the file to 0o444 after first write.""" |
| 61 | from musehub.storage.backends import LocalBackend |
| 62 | |
| 63 | with tempfile.TemporaryDirectory() as tmpdir: |
| 64 | backend = LocalBackend() |
| 65 | path = Path(tmpdir) / "repo1" / "sha256_abc" |
| 66 | path.parent.mkdir(parents=True, exist_ok=True) |
| 67 | backend._write(path, b"hello") |
| 68 | |
| 69 | mode = path.stat().st_mode |
| 70 | # No write bits for owner, group, or other. |
| 71 | assert not (mode & stat.S_IWUSR), "owner write bit must be cleared" |
| 72 | assert not (mode & stat.S_IWGRP), "group write bit must be cleared" |
| 73 | assert not (mode & stat.S_IWOTH), "other write bit must be cleared" |
| 74 | # Read bits must be set. |
| 75 | assert mode & stat.S_IRUSR |
| 76 | |
| 77 | |
| 78 | def test_local_backend_write_idempotent_for_same_content() -> None: |
| 79 | """LocalBackend._write is a no-op when the file already contains identical bytes.""" |
| 80 | from musehub.storage.backends import LocalBackend |
| 81 | |
| 82 | with tempfile.TemporaryDirectory() as tmpdir: |
| 83 | backend = LocalBackend() |
| 84 | path = Path(tmpdir) / "repo1" / "sha256_abc" |
| 85 | path.parent.mkdir(parents=True, exist_ok=True) |
| 86 | backend._write(path, b"original") |
| 87 | mtime_before = path.stat().st_mtime_ns |
| 88 | |
| 89 | # Writing the exact same bytes must not change the file. |
| 90 | backend._write(path, b"original") |
| 91 | assert path.read_bytes() == b"original", "content must be unchanged" |
| 92 | |
| 93 | |
| 94 | def test_local_backend_write_repairs_corrupt_file() -> None: |
| 95 | """LocalBackend._write overwrites when the stored bytes differ (self-heal / repair).""" |
| 96 | from musehub.storage.backends import LocalBackend |
| 97 | |
| 98 | with tempfile.TemporaryDirectory() as tmpdir: |
| 99 | backend = LocalBackend() |
| 100 | path = Path(tmpdir) / "repo1" / "sha256_abc" |
| 101 | path.parent.mkdir(parents=True, exist_ok=True) |
| 102 | path.write_bytes(b"corrupt") |
| 103 | path.chmod(0o644) |
| 104 | |
| 105 | backend._write(path, b"correct") |
| 106 | assert path.read_bytes() == b"correct", "corrupt file must be replaced" |
| 107 | |
| 108 | |
| 109 | # --------------------------------------------------------------------------- |
| 110 | # Integrity scan — clean |
| 111 | # --------------------------------------------------------------------------- |
| 112 | |
| 113 | async def test_integrity_scan_clean(db_session: AsyncSession) -> None: |
| 114 | """scan_object_integrity must return ok=True when all sampled objects are valid.""" |
| 115 | from musehub.maintenance.object_integrity import scan_object_integrity |
| 116 | |
| 117 | repo = await create_repo(db_session, slug="integrity-clean", owner="testuser") |
| 118 | data = b"test content integrity" |
| 119 | oid = blob_id(data) |
| 120 | await _add_object(db_session, repo.repo_id, data, object_id=oid) |
| 121 | await db_session.commit() |
| 122 | |
| 123 | # Backend that returns the correct content. |
| 124 | backend = AsyncMock() |
| 125 | backend.get = AsyncMock(return_value=data) |
| 126 | |
| 127 | result = await scan_object_integrity(db_session, backend, sample_size=10) |
| 128 | assert result.ok |
| 129 | assert result.mismatch_count == 0 |
| 130 | assert result.sampled >= 1 |
| 131 | |
| 132 | |
| 133 | # --------------------------------------------------------------------------- |
| 134 | # Integrity scan — missing object |
| 135 | # --------------------------------------------------------------------------- |
| 136 | |
| 137 | async def test_integrity_scan_detects_missing_object(db_session: AsyncSession) -> None: |
| 138 | """scan_object_integrity must flag objects whose backing file is absent.""" |
| 139 | from musehub.maintenance.object_integrity import scan_object_integrity |
| 140 | |
| 141 | repo = await create_repo(db_session, slug="integrity-missing", owner="testuser") |
| 142 | data = b"missing object data" |
| 143 | oid = blob_id(data) |
| 144 | await _add_object(db_session, repo.repo_id, data, object_id=oid) |
| 145 | await db_session.commit() |
| 146 | |
| 147 | backend = AsyncMock() |
| 148 | backend.get = AsyncMock(return_value=None) # simulates missing file |
| 149 | |
| 150 | result = await scan_object_integrity(db_session, backend, sample_size=10) |
| 151 | assert not result.ok |
| 152 | mismatch_ids = [m.object_id for m in result.mismatches] |
| 153 | assert oid in mismatch_ids |
| 154 | reasons = {m.reason for m in result.mismatches} |
| 155 | assert "missing" in reasons |
| 156 | |
| 157 | |
| 158 | # --------------------------------------------------------------------------- |
| 159 | # Integrity scan — hash mismatch |
| 160 | # --------------------------------------------------------------------------- |
| 161 | |
| 162 | async def test_integrity_scan_detects_hash_mismatch(db_session: AsyncSession) -> None: |
| 163 | """scan_object_integrity must flag objects whose content no longer matches their ID.""" |
| 164 | from musehub.maintenance.object_integrity import scan_object_integrity |
| 165 | |
| 166 | repo = await create_repo(db_session, slug="integrity-mismatch", owner="testuser") |
| 167 | data = b"original content" |
| 168 | oid = blob_id(data) |
| 169 | await _add_object(db_session, repo.repo_id, data, object_id=oid) |
| 170 | await db_session.commit() |
| 171 | |
| 172 | corrupted = b"corrupted content XYZ" |
| 173 | backend = AsyncMock() |
| 174 | backend.get = AsyncMock(return_value=corrupted) |
| 175 | |
| 176 | result = await scan_object_integrity(db_session, backend, sample_size=10) |
| 177 | assert not result.ok |
| 178 | mismatch_ids = [m.object_id for m in result.mismatches] |
| 179 | assert oid in mismatch_ids |
| 180 | reasons = {m.reason for m in result.mismatches} |
| 181 | assert "hash_mismatch" in reasons |
| 182 | |
| 183 | |
| 184 | # --------------------------------------------------------------------------- |
| 185 | # Soft-delete |
| 186 | # --------------------------------------------------------------------------- |
| 187 | |
| 188 | async def test_soft_delete_sets_deleted_at(db_session: AsyncSession) -> None: |
| 189 | """soft_delete_object must set deleted_at without removing the DB row.""" |
| 190 | from musehub.maintenance.object_integrity import soft_delete_object |
| 191 | from musehub.db import musehub_models as db_models |
| 192 | |
| 193 | repo = await create_repo(db_session, slug="soft-delete-test", owner="testuser") |
| 194 | data = b"soft delete me" |
| 195 | oid = blob_id(data) |
| 196 | await _add_object(db_session, repo.repo_id, data, object_id=oid) |
| 197 | await db_session.commit() |
| 198 | |
| 199 | found = await soft_delete_object(db_session, oid) |
| 200 | assert found is True |
| 201 | |
| 202 | row = await db_session.get(db_models.MusehubObject, oid) |
| 203 | assert row is not None, "row must still exist after soft delete" |
| 204 | assert row.deleted_at is not None, "deleted_at must be set" |
| 205 | |
| 206 | |
| 207 | async def test_soft_delete_unknown_object_returns_false(db_session: AsyncSession) -> None: |
| 208 | from musehub.maintenance.object_integrity import soft_delete_object |
| 209 | |
| 210 | found = await soft_delete_object(db_session, fake_id("ghost-object")) |
| 211 | assert found is False |
| 212 | |
| 213 | |
| 214 | # --------------------------------------------------------------------------- |
| 215 | # Hard-delete reaper |
| 216 | # --------------------------------------------------------------------------- |
| 217 | |
| 218 | async def test_reaper_removes_objects_past_retention(db_session: AsyncSession) -> None: |
| 219 | """reap_deleted_objects must hard-delete objects whose deleted_at exceeds retention.""" |
| 220 | from datetime import datetime, timedelta, timezone |
| 221 | from musehub.maintenance.object_integrity import reap_deleted_objects, soft_delete_object |
| 222 | from musehub.db import musehub_models as db_models |
| 223 | |
| 224 | repo = await create_repo(db_session, slug="reap-old", owner="testuser") |
| 225 | data = b"old deleted object" |
| 226 | oid = blob_id(data) |
| 227 | await _add_object(db_session, repo.repo_id, data, object_id=oid) |
| 228 | await db_session.commit() |
| 229 | |
| 230 | # Soft-delete the object, then backdate deleted_at past the retention window. |
| 231 | await soft_delete_object(db_session, oid) |
| 232 | await db_session.flush() |
| 233 | row = await db_session.get(db_models.MusehubObject, oid) |
| 234 | assert row is not None |
| 235 | row.deleted_at = datetime.now(tz=timezone.utc) - timedelta(days=31) |
| 236 | await db_session.commit() |
| 237 | |
| 238 | backend = AsyncMock() |
| 239 | backend.delete = AsyncMock() |
| 240 | |
| 241 | reaped = await reap_deleted_objects(db_session, backend, retention_days=30) |
| 242 | assert reaped == 1 |
| 243 | backend.delete.assert_called_once() |
| 244 | |
| 245 | # Row must be gone from DB. |
| 246 | gone = await db_session.get(db_models.MusehubObject, oid) |
| 247 | assert gone is None, "hard-deleted object must be removed from DB" |
| 248 | |
| 249 | |
| 250 | async def test_reaper_spares_objects_within_retention(db_session: AsyncSession) -> None: |
| 251 | """reap_deleted_objects must NOT remove objects soft-deleted within retention window.""" |
| 252 | from datetime import datetime, timedelta, timezone |
| 253 | from musehub.maintenance.object_integrity import reap_deleted_objects, soft_delete_object |
| 254 | from musehub.db import musehub_models as db_models |
| 255 | |
| 256 | repo = await create_repo(db_session, slug="reap-new", owner="testuser") |
| 257 | data = b"recently deleted object" |
| 258 | oid = blob_id(data) |
| 259 | await _add_object(db_session, repo.repo_id, data, object_id=oid) |
| 260 | await db_session.commit() |
| 261 | |
| 262 | await soft_delete_object(db_session, oid) |
| 263 | await db_session.commit() |
| 264 | |
| 265 | # deleted_at is NOW — well within the 30-day window. |
| 266 | backend = AsyncMock() |
| 267 | backend.delete = AsyncMock() |
| 268 | |
| 269 | reaped = await reap_deleted_objects(db_session, backend, retention_days=30) |
| 270 | assert reaped == 0 |
| 271 | backend.delete.assert_not_called() |
| 272 | |
| 273 | still_there = await db_session.get(db_models.MusehubObject, oid) |
| 274 | assert still_there is not None |
| 275 | |
| 276 | |
| 277 | # --------------------------------------------------------------------------- |
| 278 | # Per-repo quota — wire push |
| 279 | # --------------------------------------------------------------------------- |
| 280 | |
| 281 | async def test_wire_push_stream_rejects_when_repo_quota_exceeded(db_session: AsyncSession) -> None: |
| 282 | """wire_push_stream must reject when the push would exceed per_repo_quota_bytes.""" |
| 283 | import msgpack |
| 284 | from muse.core.mpack import MuseWireFrameWriter |
| 285 | from musehub.services.musehub_wire import wire_push_stream |
| 286 | from musehub.models.wire import SFRAME_HEADER, SFRAME_OBJECT, SFRAME_COMMIT_PACK, SFRAME_END |
| 287 | |
| 288 | repo = await create_repo(db_session, slug="quota-wire-stream", owner="testuser") |
| 289 | await db_session.commit() |
| 290 | |
| 291 | fw = MuseWireFrameWriter() |
| 292 | big_content = b"x" * 100 |
| 293 | oid = blob_id(big_content) |
| 294 | |
| 295 | def _wrap(ft: str, data: JSONValue) -> bytes: |
| 296 | return fw.wrap(frame_type=ft, payload=msgpack.packb(data, use_bin_type=True)) |
| 297 | |
| 298 | body = ( |
| 299 | _wrap(SFRAME_HEADER, {"t": SFRAME_HEADER, "branch": "main", "force": False, |
| 300 | "have": [], "head": fake_id("push-head"), "n_objects": 1, "n_commits": 0}) |
| 301 | + _wrap(SFRAME_OBJECT, {"t": SFRAME_OBJECT, "id": oid, "content": big_content, "path": "big.bin", "enc": "raw"}) |
| 302 | + _wrap(SFRAME_COMMIT_PACK, {"t": SFRAME_COMMIT_PACK, "commits": [], "snapshots": []}) |
| 303 | + _wrap(SFRAME_END, {"t": SFRAME_END, "n_objects": 1, "n_commits": 0}) |
| 304 | ) |
| 305 | |
| 306 | async def body_iter() -> None: |
| 307 | yield body |
| 308 | |
| 309 | frames: list[dict] = [] |
| 310 | with patch("musehub.services.musehub_wire.settings") as mock_settings: |
| 311 | mock_settings.per_repo_quota_bytes = 50 # 50 bytes — 100-byte object exceeds it |
| 312 | mock_settings.require_signed_commits = False |
| 313 | mock_settings.trusted_agent_ids = [] |
| 314 | async for chunk in wire_push_stream(db_session, repo.repo_id, body_iter(), pusher_id="testuser"): |
| 315 | unpacker = msgpack.Unpacker(raw=False) |
| 316 | unpacker.feed(chunk) |
| 317 | frames.extend(list(unpacker)) |
| 318 | |
| 319 | result = frames[-1] if frames else {} |
| 320 | assert result.get("ok") is not True # X frame (error) has no "ok"; R frame has ok=False |
| 321 | assert "quota" in (result.get("msg") or result.get("message") or "").lower() |
File History
1 commit
sha256:a34090cc4a394a78bd72cbbe34b08cc59525141e19135b6c0ab154f10611b9ef
debug(push/stream): instrument O-frame decode path with INF…
Sonnet 4.6
patch
121 days ago