"""Tests for checklist section 5.2 — Object store. Covers: - Object files are immutable after write (LocalBackend chmod 0o444) - Integrity scan detects missing objects - Integrity scan detects hash-mismatched objects - Integrity scan passes on clean objects - Per-repo storage quota enforced in wire push - Soft-delete sets deleted_at without removing the file - Hard-delete reaper removes only objects past retention window """ from __future__ import annotations import stat import tempfile from pathlib import Path from unittest.mock import AsyncMock, MagicMock, patch import pytest from sqlalchemy.ext.asyncio import AsyncSession from muse.core.types import blob_id, fake_id from musehub.db.musehub_models import MusehubObject as _MusehubObject from musehub.types.json_types import JSONValue from tests.factories import create_repo # --------------------------------------------------------------------------- # Helpers # --------------------------------------------------------------------------- async def _add_object( session: AsyncSession, repo_id: str, data: bytes, *, object_id: str | None = None, ) -> _MusehubObject: from musehub.db import musehub_models as db oid = object_id or blob_id(data) obj = db.MusehubObject( object_id=oid, path="test.bin", size_bytes=len(data), disk_path=f"/tmp/{oid.replace(':', '_')}", storage_uri=f"local:///tmp/{oid.replace(':', '_')}", ) session.add(obj) session.add(db.MusehubObjectRef(repo_id=repo_id, object_id=oid)) await session.flush() return obj # --------------------------------------------------------------------------- # Immutability — LocalBackend sets file to read-only after write # --------------------------------------------------------------------------- def test_local_backend_write_sets_readonly() -> None: """LocalBackend._write must chmod the file to 0o444 after first write.""" from musehub.storage.backends import LocalBackend with tempfile.TemporaryDirectory() as tmpdir: backend = LocalBackend() path = Path(tmpdir) / "repo1" / "sha256_abc" path.parent.mkdir(parents=True, exist_ok=True) backend._write(path, b"hello") mode = path.stat().st_mode # No write bits for owner, group, or other. assert not (mode & stat.S_IWUSR), "owner write bit must be cleared" assert not (mode & stat.S_IWGRP), "group write bit must be cleared" assert not (mode & stat.S_IWOTH), "other write bit must be cleared" # Read bits must be set. assert mode & stat.S_IRUSR def test_local_backend_write_idempotent_for_same_content() -> None: """LocalBackend._write is a no-op when the file already contains identical bytes.""" from musehub.storage.backends import LocalBackend with tempfile.TemporaryDirectory() as tmpdir: backend = LocalBackend() path = Path(tmpdir) / "repo1" / "sha256_abc" path.parent.mkdir(parents=True, exist_ok=True) backend._write(path, b"original") mtime_before = path.stat().st_mtime_ns # Writing the exact same bytes must not change the file. backend._write(path, b"original") assert path.read_bytes() == b"original", "content must be unchanged" def test_local_backend_write_repairs_corrupt_file() -> None: """LocalBackend._write overwrites when the stored bytes differ (self-heal / repair).""" from musehub.storage.backends import LocalBackend with tempfile.TemporaryDirectory() as tmpdir: backend = LocalBackend() path = Path(tmpdir) / "repo1" / "sha256_abc" path.parent.mkdir(parents=True, exist_ok=True) path.write_bytes(b"corrupt") path.chmod(0o644) backend._write(path, b"correct") assert path.read_bytes() == b"correct", "corrupt file must be replaced" # --------------------------------------------------------------------------- # Integrity scan — clean # --------------------------------------------------------------------------- async def test_integrity_scan_clean(db_session: AsyncSession) -> None: """scan_object_integrity must return ok=True when all sampled objects are valid.""" from musehub.maintenance.object_integrity import scan_object_integrity repo = await create_repo(db_session, slug="integrity-clean", owner="testuser") data = b"test content integrity" oid = blob_id(data) await _add_object(db_session, repo.repo_id, data, object_id=oid) await db_session.commit() # Backend that returns the correct content. backend = AsyncMock() backend.get = AsyncMock(return_value=data) result = await scan_object_integrity(db_session, backend, sample_size=10) assert result.ok assert result.mismatch_count == 0 assert result.sampled >= 1 # --------------------------------------------------------------------------- # Integrity scan — missing object # --------------------------------------------------------------------------- async def test_integrity_scan_detects_missing_object(db_session: AsyncSession) -> None: """scan_object_integrity must flag objects whose backing file is absent.""" from musehub.maintenance.object_integrity import scan_object_integrity repo = await create_repo(db_session, slug="integrity-missing", owner="testuser") data = b"missing object data" oid = blob_id(data) await _add_object(db_session, repo.repo_id, data, object_id=oid) await db_session.commit() backend = AsyncMock() backend.get = AsyncMock(return_value=None) # simulates missing file result = await scan_object_integrity(db_session, backend, sample_size=10) assert not result.ok mismatch_ids = [m.object_id for m in result.mismatches] assert oid in mismatch_ids reasons = {m.reason for m in result.mismatches} assert "missing" in reasons # --------------------------------------------------------------------------- # Integrity scan — hash mismatch # --------------------------------------------------------------------------- async def test_integrity_scan_detects_hash_mismatch(db_session: AsyncSession) -> None: """scan_object_integrity must flag objects whose content no longer matches their ID.""" from musehub.maintenance.object_integrity import scan_object_integrity repo = await create_repo(db_session, slug="integrity-mismatch", owner="testuser") data = b"original content" oid = blob_id(data) await _add_object(db_session, repo.repo_id, data, object_id=oid) await db_session.commit() corrupted = b"corrupted content XYZ" backend = AsyncMock() backend.get = AsyncMock(return_value=corrupted) result = await scan_object_integrity(db_session, backend, sample_size=10) assert not result.ok mismatch_ids = [m.object_id for m in result.mismatches] assert oid in mismatch_ids reasons = {m.reason for m in result.mismatches} assert "hash_mismatch" in reasons # --------------------------------------------------------------------------- # Soft-delete # --------------------------------------------------------------------------- async def test_soft_delete_sets_deleted_at(db_session: AsyncSession) -> None: """soft_delete_object must set deleted_at without removing the DB row.""" from musehub.maintenance.object_integrity import soft_delete_object from musehub.db import musehub_models as db_models repo = await create_repo(db_session, slug="soft-delete-test", owner="testuser") data = b"soft delete me" oid = blob_id(data) await _add_object(db_session, repo.repo_id, data, object_id=oid) await db_session.commit() found = await soft_delete_object(db_session, oid) assert found is True row = await db_session.get(db_models.MusehubObject, oid) assert row is not None, "row must still exist after soft delete" assert row.deleted_at is not None, "deleted_at must be set" async def test_soft_delete_unknown_object_returns_false(db_session: AsyncSession) -> None: from musehub.maintenance.object_integrity import soft_delete_object found = await soft_delete_object(db_session, fake_id("ghost-object")) assert found is False # --------------------------------------------------------------------------- # Hard-delete reaper # --------------------------------------------------------------------------- async def test_reaper_removes_objects_past_retention(db_session: AsyncSession) -> None: """reap_deleted_objects must hard-delete objects whose deleted_at exceeds retention.""" from datetime import datetime, timedelta, timezone from musehub.maintenance.object_integrity import reap_deleted_objects, soft_delete_object from musehub.db import musehub_models as db_models repo = await create_repo(db_session, slug="reap-old", owner="testuser") data = b"old deleted object" oid = blob_id(data) await _add_object(db_session, repo.repo_id, data, object_id=oid) await db_session.commit() # Soft-delete the object, then backdate deleted_at past the retention window. await soft_delete_object(db_session, oid) await db_session.flush() row = await db_session.get(db_models.MusehubObject, oid) assert row is not None row.deleted_at = datetime.now(tz=timezone.utc) - timedelta(days=31) await db_session.commit() backend = AsyncMock() backend.delete = AsyncMock() reaped = await reap_deleted_objects(db_session, backend, retention_days=30) assert reaped == 1 backend.delete.assert_called_once() # Row must be gone from DB. gone = await db_session.get(db_models.MusehubObject, oid) assert gone is None, "hard-deleted object must be removed from DB" async def test_reaper_spares_objects_within_retention(db_session: AsyncSession) -> None: """reap_deleted_objects must NOT remove objects soft-deleted within retention window.""" from datetime import datetime, timedelta, timezone from musehub.maintenance.object_integrity import reap_deleted_objects, soft_delete_object from musehub.db import musehub_models as db_models repo = await create_repo(db_session, slug="reap-new", owner="testuser") data = b"recently deleted object" oid = blob_id(data) await _add_object(db_session, repo.repo_id, data, object_id=oid) await db_session.commit() await soft_delete_object(db_session, oid) await db_session.commit() # deleted_at is NOW — well within the 30-day window. backend = AsyncMock() backend.delete = AsyncMock() reaped = await reap_deleted_objects(db_session, backend, retention_days=30) assert reaped == 0 backend.delete.assert_not_called() still_there = await db_session.get(db_models.MusehubObject, oid) assert still_there is not None # --------------------------------------------------------------------------- # Per-repo quota — wire push # --------------------------------------------------------------------------- async def test_wire_push_stream_rejects_when_repo_quota_exceeded(db_session: AsyncSession) -> None: """wire_push_stream must reject when the push would exceed per_repo_quota_bytes.""" import msgpack from muse.core.mpack import MuseWireFrameWriter from musehub.services.musehub_wire import wire_push_stream from musehub.models.wire import SFRAME_HEADER, SFRAME_OBJECT, SFRAME_COMMIT_PACK, SFRAME_END repo = await create_repo(db_session, slug="quota-wire-stream", owner="testuser") await db_session.commit() fw = MuseWireFrameWriter() big_content = b"x" * 100 oid = blob_id(big_content) def _wrap(ft: str, data: JSONValue) -> bytes: return fw.wrap(frame_type=ft, payload=msgpack.packb(data, use_bin_type=True)) body = ( _wrap(SFRAME_HEADER, {"t": SFRAME_HEADER, "branch": "main", "force": False, "have": [], "head": fake_id("push-head"), "n_objects": 1, "n_commits": 0}) + _wrap(SFRAME_OBJECT, {"t": SFRAME_OBJECT, "id": oid, "content": big_content, "path": "big.bin", "enc": "raw"}) + _wrap(SFRAME_COMMIT_PACK, {"t": SFRAME_COMMIT_PACK, "commits": [], "snapshots": []}) + _wrap(SFRAME_END, {"t": SFRAME_END, "n_objects": 1, "n_commits": 0}) ) async def body_iter() -> None: yield body frames: list[dict] = [] with patch("musehub.services.musehub_wire.settings") as mock_settings: mock_settings.per_repo_quota_bytes = 50 # 50 bytes — 100-byte object exceeds it mock_settings.require_signed_commits = False mock_settings.trusted_agent_ids = [] async for chunk in wire_push_stream(db_session, repo.repo_id, body_iter(), pusher_id="testuser"): unpacker = msgpack.Unpacker(raw=False) unpacker.feed(chunk) frames.extend(list(unpacker)) result = frames[-1] if frames else {} assert result.get("ok") is not True # X frame (error) has no "ok"; R frame has ok=False assert "quota" in (result.get("msg") or result.get("message") or "").lower()