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