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