gabriel / musehub public
test_mpack_index_phase1.py python
241 lines 8.8 KB
Raw
sha256:ef10830ce231e0a20efcb0e2586cb879471247e916616e6fdd0d51df459e2595 fix: typing audit — 0 violations, 0 untyped defs across all… Sonnet 4.6 minor ⚠ breaking 58 days ago
1 """TDD — Phase 1: mpack index written for every pushed object (issue #63).
2
3 PI-1 After process_mpack_index_job, every object in the mpack has a row
4 in musehub_mpack_index with the correct mpack_id (= mpack_key).
5 PI-2 The mpack_id stored matches the mpack_key exactly.
6 PI-3 A second push of different objects adds new rows; existing rows are
7 untouched (on_conflict_do_nothing).
8 PI-4 Objects from different repos are distinct in the global index
9 (content-addressed objects are globally unique by object_id).
10 """
11 from __future__ import annotations
12
13 import datetime
14 import hashlib
15 from collections.abc import Mapping
16
17 import msgpack
18 import pytest
19 from sqlalchemy import select
20 from sqlalchemy.ext.asyncio import AsyncSession
21
22 from muse.core.types import blob_id
23 from musehub.db import musehub_repo_models as db
24 from musehub.core.genesis import compute_identity_id
25 from musehub.services.musehub_repository import create_repo
26
27
28 # ---------------------------------------------------------------------------
29 # Helpers
30 # ---------------------------------------------------------------------------
31
32 def _make_mpack(objects: Mapping[str, bytes]) -> tuple[bytes, str]:
33 """Build a minimal MPack and return (wire_bytes, mpack_key)."""
34 mpack = {
35 "commits": [],
36 "snapshots": [],
37 "objects": [
38 {"object_id": oid, "content": data}
39 for oid, data in objects.items()
40 ],
41 "branch_heads": {},
42 }
43 wire_bytes = msgpack.packb(mpack, use_bin_type=True)
44 mpack_key = "sha256:" + hashlib.sha256(wire_bytes).hexdigest()
45 return wire_bytes, mpack_key
46
47
48 async def _store_mpack(wire_bytes: bytes, mpack_key: str) -> None:
49 """Write mpack bytes to storage so process_mpack_index_job can fetch it."""
50 import musehub.storage.backends as _backends_mod
51 backend = _backends_mod.get_backend()
52 await backend.put_mpack(mpack_key, wire_bytes)
53
54
55 async def _enqueue_and_process(
56 session: AsyncSession,
57 repo_id: str,
58 mpack_key: str,
59 n_objects: int,
60 ) -> str:
61 """Enqueue a mpack.index job and run it synchronously. Returns job_id."""
62 from musehub.core.genesis import compute_job_id
63 from musehub.db.musehub_jobs_models import MusehubBackgroundJob
64 from musehub.services.musehub_wire import process_mpack_index_job
65
66 now = datetime.datetime.now(datetime.timezone.utc)
67 job_id = compute_job_id(repo_id, "mpack.index", now.isoformat())
68 session.add(MusehubBackgroundJob(
69 job_id=job_id,
70 repo_id=repo_id,
71 job_type="mpack.index",
72 payload={
73 "mpack_key": mpack_key,
74 "branch": "main",
75 "head": "",
76 "pusher_id": "",
77 "declared_objects_count": n_objects,
78 "declared_commits_count": 0,
79 },
80 status="pending",
81 created_at=now,
82 attempt=0,
83 ))
84 await session.commit()
85 await process_mpack_index_job(session, job_id)
86 await session.commit()
87 return job_id
88
89
90 # ---------------------------------------------------------------------------
91 # PI-1
92 # ---------------------------------------------------------------------------
93
94 @pytest.mark.asyncio
95 async def test_pi1_mpack_index_written_for_every_object(db_session: AsyncSession) -> None:
96 """Every object in the mpack must have a row in musehub_mpack_index."""
97 repo = await create_repo(
98 db_session,
99 name="pi-test-1",
100 owner="gabriel",
101 owner_user_id=compute_identity_id(b"gabriel"),
102 visibility="public",
103 initialize=False,
104 )
105
106 objects = {blob_id(f"pi1-obj-{i}".encode()): f"pi1-obj-{i}".encode() for i in range(5)}
107 wire_bytes, mpack_key = _make_mpack(objects)
108 await _store_mpack(wire_bytes, mpack_key)
109 await _enqueue_and_process(db_session, repo.repo_id, mpack_key, len(objects))
110
111 rows_q = await db_session.execute(
112 select(db.MusehubMPackIndex).where(
113 db.MusehubMPackIndex.entity_id.in_(list(objects.keys()))
114 )
115 )
116 rows = rows_q.scalars().all()
117 indexed_oids = {r.entity_id for r in rows}
118
119 assert indexed_oids == set(objects.keys()), (
120 f"expected {len(objects)} mpack index rows, got {len(rows)}\n"
121 f"missing: {set(objects.keys()) - indexed_oids}"
122 )
123
124
125 # ---------------------------------------------------------------------------
126 # PI-2
127 # ---------------------------------------------------------------------------
128
129 @pytest.mark.asyncio
130 async def test_pi2_mpack_id_matches_mpack_key(db_session: AsyncSession) -> None:
131 """The mpack_id stored in musehub_mpack_index must equal the mpack_key."""
132 repo = await create_repo(
133 db_session,
134 name="pi-test-2",
135 owner="gabriel",
136 owner_user_id=compute_identity_id(b"gabriel"),
137 visibility="public",
138 initialize=False,
139 )
140
141 objects = {blob_id(b"pi2-only-obj"): b"pi2-only-obj"}
142 wire_bytes, mpack_key = _make_mpack(objects)
143 await _store_mpack(wire_bytes, mpack_key)
144 await _enqueue_and_process(db_session, repo.repo_id, mpack_key, len(objects))
145
146 row = await db_session.get(db.MusehubMPackIndex, (list(objects.keys())[0], mpack_key))
147 assert row is not None, "mpack index row not found"
148 assert row.mpack_id == mpack_key, f"mpack_id {row.mpack_id!r} != mpack_key {mpack_key!r}"
149
150
151 # ---------------------------------------------------------------------------
152 # PI-3
153 # ---------------------------------------------------------------------------
154
155 @pytest.mark.asyncio
156 async def test_pi3_second_push_adds_rows_without_overwriting(db_session: AsyncSession) -> None:
157 """A second push of different objects adds new rows; first push rows survive."""
158 repo = await create_repo(
159 db_session,
160 name="pi-test-3",
161 owner="gabriel",
162 owner_user_id=compute_identity_id(b"gabriel"),
163 visibility="public",
164 initialize=False,
165 )
166
167 objs_a = {blob_id(f"pi3-a-{i}".encode()): f"pi3-a-{i}".encode() for i in range(3)}
168 wire_a, key_a = _make_mpack(objs_a)
169 await _store_mpack(wire_a, key_a)
170 await _enqueue_and_process(db_session, repo.repo_id, key_a, len(objs_a))
171
172 objs_b = {blob_id(f"pi3-b-{i}".encode()): f"pi3-b-{i}".encode() for i in range(3)}
173 wire_b, key_b = _make_mpack(objs_b)
174 await _store_mpack(wire_b, key_b)
175 await _enqueue_and_process(db_session, repo.repo_id, key_b, len(objs_b))
176
177 all_oids = set(objs_a.keys()) | set(objs_b.keys())
178 rows_q = await db_session.execute(
179 select(db.MusehubMPackIndex).where(db.MusehubMPackIndex.entity_id.in_(list(all_oids)))
180 )
181 rows = rows_q.scalars().all()
182 indexed_oids = {r.entity_id for r in rows}
183 indexed_mpacks = {r.mpack_id for r in rows}
184
185 assert set(objs_a.keys()) <= indexed_oids, "first push objects missing from index"
186 assert set(objs_b.keys()) <= indexed_oids, "second push objects missing from index"
187 assert key_a in indexed_mpacks, "first mpack_key missing from index"
188 assert key_b in indexed_mpacks, "second mpack_key missing from index"
189
190
191 # ---------------------------------------------------------------------------
192 # PI-4
193 # ---------------------------------------------------------------------------
194
195 @pytest.mark.asyncio
196 async def test_pi4_objects_from_different_repos_are_distinct(db_session: AsyncSession) -> None:
197 """Objects pushed to different repos are distinct in the global index (no collisions)."""
198 repo_a = await create_repo(
199 db_session,
200 name="pi-test-4a",
201 owner="gabriel",
202 owner_user_id=compute_identity_id(b"gabriel"),
203 visibility="public",
204 initialize=False,
205 )
206 repo_b = await create_repo(
207 db_session,
208 name="pi-test-4b",
209 owner="gabriel",
210 owner_user_id=compute_identity_id(b"gabriel"),
211 visibility="public",
212 initialize=False,
213 )
214
215 objs_a = {blob_id(b"pi4-repo-a-obj"): b"pi4-repo-a-obj"}
216 wire_a, key_a = _make_mpack(objs_a)
217 await _store_mpack(wire_a, key_a)
218 await _enqueue_and_process(db_session, repo_a.repo_id, key_a, 1)
219
220 objs_b = {blob_id(b"pi4-repo-b-obj"): b"pi4-repo-b-obj"}
221 wire_b, key_b = _make_mpack(objs_b)
222 await _store_mpack(wire_b, key_b)
223 await _enqueue_and_process(db_session, repo_b.repo_id, key_b, 1)
224
225 rows_a = (await db_session.execute(
226 select(db.MusehubMPackIndex).where(
227 db.MusehubMPackIndex.entity_id.in_(list(objs_a.keys()))
228 )
229 )).scalars().all()
230
231 rows_b = (await db_session.execute(
232 select(db.MusehubMPackIndex).where(
233 db.MusehubMPackIndex.entity_id.in_(list(objs_b.keys()))
234 )
235 )).scalars().all()
236
237 assert {r.entity_id for r in rows_a} == set(objs_a.keys()), "repo_a objects not indexed"
238 assert {r.entity_id for r in rows_b} == set(objs_b.keys()), "repo_b objects not indexed"
239 assert {r.entity_id for r in rows_a}.isdisjoint({r.entity_id for r in rows_b}), (
240 "object collision — two repos pushed the same object_id (seeds must be unique)"
241 )
File History 1 commit
sha256:ef10830ce231e0a20efcb0e2586cb879471247e916616e6fdd0d51df459e2595 fix: typing audit — 0 violations, 0 untyped defs across all… Sonnet 4.6 minor 58 days ago