gabriel / musehub public
test_object_store.py python
325 lines 12.9 KB
Raw
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