"""TDD — CF Worker pack-receiver routing (Phase D). Tests that FAIL before implementation (functions don't exist yet): TestPostObjectPack::* — _post_object_pack not yet in push.py TestPushStreamWorkerRouting::* — _push_stream doesn't accept pack_origin yet Tests that PASS before implementation (already wired on server side): TestParseRemoteInfoPackOrigin::* — _parse_remote_info already reads pack_origin After implementation all tests must pass. """ from __future__ import annotations import pathlib import msgpack import pytest from unittest.mock import AsyncMock, MagicMock, patch from muse.core.pack import ObjectPayload, RemoteInfo from muse.core._types import MsgpackDict from muse.core.transport import _parse_remote_info # ── RemoteInfo / refs parsing ───────────────────────────────────────────────── class TestParseRemoteInfoPackOrigin: def _raw(self, **extra) -> bytes: return msgpack.packb( {"repo_id": "r1", "domain": "code", "default_branch": "main", "branch_heads": {}, **extra}, use_bin_type=True, ) def test_pack_origin_parsed(self) -> None: info = _parse_remote_info(self._raw(pack_origin="https://pack.staging.musehub.ai")) assert info.get("pack_origin") == "https://pack.staging.musehub.ai" def test_pack_origin_absent_returns_none(self) -> None: info = _parse_remote_info(self._raw()) assert info.get("pack_origin") is None def test_pack_origin_whitespace_only_ignored(self) -> None: info = _parse_remote_info(self._raw(pack_origin=" ")) assert info.get("pack_origin") is None def test_pack_origin_stored_in_remote_info(self) -> None: info = _parse_remote_info(self._raw(pack_origin="https://pack.musehub.ai")) assert "pack_origin" in info # ── _post_object_pack unit tests ────────────────────────────────────────────── class TestPostObjectPack: """Unit tests for the Worker upload coroutine.""" def _make_objects(self, n: int = 2) -> list[ObjectPayload]: return [ ObjectPayload( object_id="sha256:" + hex(i)[2:].zfill(64), content=f"content{i}".encode(), path=f"track{i}.mid", ) for i in range(n) ] def _mock_client(self, status: int = 200, body: MsgpackDict | None = None) -> AsyncMock: mock_resp = MagicMock() mock_resp.status_code = status mock_resp.headers = {"content-type": "application/json"} mock_resp.json.return_value = body or {"stored": len(self._make_objects()), "skipped": 0} mock_resp.text = "error" if status >= 400 else "" client = AsyncMock() client.post.return_value = mock_resp return client @pytest.mark.asyncio async def test_posts_to_provided_url(self) -> None: from muse.cli.commands.push import _post_object_pack client = self._mock_client(body={"stored": 1, "skipped": 0}) await _post_object_pack( "https://pack.musehub.ai/gabriel/my-repo/push/object-pack", self._make_objects(1), None, client, ) client.post.assert_called_once() assert client.post.call_args[0][0] == ( "https://pack.musehub.ai/gabriel/my-repo/push/object-pack" ) @pytest.mark.asyncio async def test_body_is_msgpack_with_objects_key(self) -> None: from muse.cli.commands.push import _post_object_pack objs = self._make_objects(2) client = self._mock_client(body={"stored": 2, "skipped": 0}) await _post_object_pack("https://pack.example.com/a/b/push/object-pack", objs, None, client) sent_body: bytes = client.post.call_args.kwargs["content"] decoded = msgpack.unpackb(sent_body, raw=False) assert "objects" in decoded assert len(decoded["objects"]) == 2 assert decoded["objects"][0]["object_id"] == objs[0]["object_id"] assert decoded["objects"][0]["content"] == b"content0" assert decoded["objects"][0]["path"] == "track0.mid" @pytest.mark.asyncio async def test_returns_stored_skipped_counts(self) -> None: from muse.cli.commands.push import _post_object_pack client = self._mock_client(body={"stored": 3, "skipped": 1}) result = await _post_object_pack( "https://pack.example.com/a/b/push/object-pack", [], None, client, ) assert result == {"stored": 3, "skipped": 1} @pytest.mark.asyncio async def test_accepts_msgpack_response(self) -> None: from muse.cli.commands.push import _post_object_pack client = AsyncMock() mock_resp = MagicMock() mock_resp.status_code = 200 mock_resp.headers = {"content-type": "application/x-msgpack"} mock_resp.content = msgpack.packb({"stored": 5, "skipped": 2}, use_bin_type=True) client.post.return_value = mock_resp result = await _post_object_pack("https://pack.example.com/x/y/push/object-pack", [], None, client) assert result == {"stored": 5, "skipped": 2} @pytest.mark.asyncio async def test_raises_transport_error_on_4xx(self) -> None: from muse.cli.commands.push import _post_object_pack from muse.core.transport import TransportError client = self._mock_client(status=403) with pytest.raises(TransportError) as exc_info: await _post_object_pack("https://pack.example.com/x/y/push/object-pack", [], None, client) assert exc_info.value.status_code == 403 @pytest.mark.asyncio async def test_raises_transport_error_on_5xx(self) -> None: from muse.cli.commands.push import _post_object_pack from muse.core.transport import TransportError client = self._mock_client(status=502) with pytest.raises(TransportError) as exc_info: await _post_object_pack("https://pack.example.com/x/y/push/object-pack", [], None, client) assert exc_info.value.status_code == 502 @pytest.mark.asyncio async def test_content_type_header_is_msgpack(self) -> None: from muse.cli.commands.push import _post_object_pack client = self._mock_client(body={"stored": 0, "skipped": 0}) await _post_object_pack("https://pack.example.com/a/b/push/object-pack", [], None, client) sent_headers = client.post.call_args.kwargs.get("headers", {}) assert sent_headers.get("Content-Type") == "application/x-msgpack" # ── _push_stream Worker routing integration ─────────────────────────────────── def _fresh_repo(tmp_path: pathlib.Path, monkeypatch: pytest.MonkeyPatch) -> pathlib.Path: """Create a repo with one commit and one object.""" from tests.cli_test_helper import CliRunner runner = CliRunner() cli = None env = {"MUSE_REPO_ROOT": str(tmp_path)} monkeypatch.chdir(tmp_path) monkeypatch.setenv("MUSE_REPO_ROOT", str(tmp_path)) r = runner.invoke(cli, ["init"], env=env, catch_exceptions=False) assert r.exit_code == 0, r.output (tmp_path / "track.mid").write_bytes(b"\x00MIDI" + b"\xff" * 200) r = runner.invoke(cli, ["code", "add", "track.mid"], env=env, catch_exceptions=False) assert r.exit_code == 0, r.output r = runner.invoke(cli, ["commit", "-m", "add track"], env=env, catch_exceptions=False) assert r.exit_code == 0, r.output return tmp_path class TestPushStreamWorkerRouting: """_push_stream routes correctly based on whether pack_origin is set.""" def _fake_transport(self) -> MagicMock: transport = MagicMock() transport.fetch_remote_info.return_value = RemoteInfo( repo_id="r1", domain="midi", default_branch="main", branch_heads={}, ) from muse.core.pack import PushResult ok_result = PushResult(ok=True, message="pushed", branch_heads={"main": "sha256:" + "a" * 64}) coro = AsyncMock(return_value=ok_result) transport.push_stream_coro = coro return transport def test_push_stream_accepts_pack_origin_param(self, tmp_path: pathlib.Path) -> None: """_push_stream signature must accept pack_origin keyword argument.""" import inspect from muse.cli.commands.push import _push_stream sig = inspect.signature(_push_stream) assert "pack_origin" in sig.parameters, ( "_push_stream must accept pack_origin parameter" ) def test_objects_sent_to_worker_not_stream_when_pack_origin_set( self, tmp_path: pathlib.Path, monkeypatch: pytest.MonkeyPatch, ) -> None: """When pack_origin is set, objects go to Worker; push_stream_coro gets objects=[].""" repo = _fresh_repo(tmp_path, monkeypatch) from muse.core.store import get_head_commit_id local_head = get_head_commit_id(repo, "main") assert local_head transport = self._fake_transport() worker_calls: list[tuple] = [] async def fake_post_object_pack(worker_url, objects, signing, client): worker_calls.append((worker_url, list(objects))) return {"stored": len(objects), "skipped": 0} from muse.cli.commands.push import _push_stream with patch("muse.cli.commands.push._post_object_pack", fake_post_object_pack): result, _commits, _objects = _push_stream( transport, url="https://staging.musehub.ai/gabriel/test-repo", signing=None, root=repo, local_head=local_head, have=[], branch="main", force=False, pack_origin="https://pack.staging.musehub.ai", ) # Worker must have been called with objects assert len(worker_calls) >= 1, "Worker was never called" all_worker_objects = [o for _, objs in worker_calls for o in objs] assert len(all_worker_objects) >= 1, "No objects sent to Worker" # push_stream_coro must have been called with objects=[] stream_call = transport.push_stream_coro.call_args assert stream_call is not None, "push_stream_coro was never called" assert stream_call.kwargs.get("objects") == [], ( f"Expected objects=[] in stream call, got {stream_call.kwargs.get('objects')}" ) def test_worker_url_constructed_from_pack_origin_and_repo_path( self, tmp_path: pathlib.Path, monkeypatch: pytest.MonkeyPatch, ) -> None: """Worker URL = {pack_origin}{/owner/slug}/push/object-pack.""" repo = _fresh_repo(tmp_path, monkeypatch) from muse.core.store import get_head_commit_id local_head = get_head_commit_id(repo, "main") transport = self._fake_transport() captured_urls: list[str] = [] async def fake_post_object_pack(worker_url, objects, signing, client): captured_urls.append(worker_url) return {"stored": len(objects), "skipped": 0} from muse.cli.commands.push import _push_stream with patch("muse.cli.commands.push._post_object_pack", fake_post_object_pack): _push_stream( transport, url="https://staging.musehub.ai/gabriel/test-repo", signing=None, root=repo, local_head=local_head, have=[], branch="main", force=False, pack_origin="https://pack.staging.musehub.ai", ) assert len(captured_urls) >= 1 assert captured_urls[0] == ( "https://pack.staging.musehub.ai/gabriel/test-repo/push/object-pack" ) def test_no_worker_call_when_pack_origin_absent( self, tmp_path: pathlib.Path, monkeypatch: pytest.MonkeyPatch, ) -> None: """Without pack_origin, _post_object_pack is never called.""" repo = _fresh_repo(tmp_path, monkeypatch) from muse.core.store import get_head_commit_id local_head = get_head_commit_id(repo, "main") transport = self._fake_transport() worker_calls: list[str] = [] async def fake_post_object_pack(worker_url, objects, signing, client): worker_calls.append(worker_url) return {"stored": 0, "skipped": 0} from muse.cli.commands.push import _push_stream with patch("muse.cli.commands.push._post_object_pack", fake_post_object_pack): _push_stream( transport, url="https://staging.musehub.ai/gabriel/test-repo", signing=None, root=repo, local_head=local_head, have=[], branch="main", force=False, pack_origin=None, # no Worker ) assert worker_calls == [], "Worker must not be called when pack_origin is absent" def test_stream_carries_objects_when_no_pack_origin( self, tmp_path: pathlib.Path, monkeypatch: pytest.MonkeyPatch, ) -> None: """Without pack_origin, push_stream_coro receives the objects.""" repo = _fresh_repo(tmp_path, monkeypatch) from muse.core.store import get_head_commit_id local_head = get_head_commit_id(repo, "main") transport = self._fake_transport() from muse.cli.commands.push import _push_stream _push_stream( transport, url="https://staging.musehub.ai/gabriel/test-repo", signing=None, root=repo, local_head=local_head, have=[], branch="main", force=False, pack_origin=None, ) stream_call = transport.push_stream_coro.call_args assert stream_call is not None stream_objects = stream_call.kwargs.get("objects", []) assert len(stream_objects) >= 1, ( "Without pack_origin, objects must travel via push/stream" )