test_mpack_index_phase1.py
python
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