"""Round 68: the request / delivery loop. * Sender: their requests ∩ what I offer *now* → my outbox (copied under a .part name Syncthing never sees, then renamed); anything no longer requested is deleted — that's how the outbox cleans itself up. * Receiver: a finished file I want is imported like Add to Library (date added now, their listening history left behind) or, for a followed playlist, copied to the cache; then my requests are rewritten without it and without anything they stopped offering. * Safe to repeat: Syncthing temp files are ignored, an exact match already in my library is never imported twice. ``test_loop_between_two_libraries`` stands in for Syncthing by copying each side's out/ into the other's in/; the ``syncthing``-marked test does it for real. """ import shutil import time from pathlib import Path import pytest from lintunes.cassette import delivery, matching, publish, share from lintunes.cassette import requests as req from lintunes.cassette.service import CassetteService from lintunes.cassette.state import ( FOLLOW_LIBRARY, Friend, Selection, cache_dir, in_dir, out_dir, ) from lintunes.models import Playlist, Track from lintunes.models.library import Library def _library_with_file(tmp_path, mp3_file, tid=5, name="Song", artist="Artist", album="Album"): library = Library() path = tmp_path / f"{name}.mp3" shutil.copyfile(mp3_file, path) library.tracks[tid] = Track(track_id=tid, name=name, artist=artist, album=album, location=str(path), size=path.stat().st_size, total_time=1000, play_count=99, rating=100) return library # ---- names ---- class TestNames: def test_outbox_names(self): name = delivery.outbox_name(12, "/m/A/B/04 Hey: Jude.mp3") assert name == "12 ~ 04 Hey_ Jude.mp3" assert delivery.parse_outbox_name(name) == (12, "04 Hey_ Jude.mp3") assert delivery.parse_outbox_name("notes.txt") is None @pytest.mark.parametrize("name", [ ".syncthing.12 ~ a.mp3.tmp", "~syncthing~12 ~ a.mp3.tmp", ".12 ~ a.mp3.part"]) def test_temp_files_are_invisible(self, tmp_path, name): (tmp_path / name).write_text("x") assert delivery.list_outbox(tmp_path) == {} # ---- sending ---- class TestSend: def test_requested_and_offered_only(self, tmp_path, mp3_file): library = _library_with_file(tmp_path, mp3_file) library.tracks[6] = Track(track_id=6, name="Not offered", location=library.tracks[5].location) root = tmp_path / "root" req.write_requests(in_dir(root, "t"), [5, 6]) plan = delivery.plan_send("t", root, library, Selection(track_ids={5})) assert [name for _src, name in plan.copies] == ["5 ~ Song.mp3"] def test_stopped_offering_means_not_sent(self, tmp_path, mp3_file): library = _library_with_file(tmp_path, mp3_file) root = tmp_path / "root" req.write_requests(in_dir(root, "t"), [5]) plan = delivery.plan_send("t", root, library, Selection()) assert plan.copies == [] def test_copy_then_clean_up(self, tmp_path, mp3_file): library = _library_with_file(tmp_path, mp3_file) root = tmp_path / "root" selection = Selection(all_library=True) req.write_requests(in_dir(root, "t"), [5]) plan = delivery.plan_send("t", root, library, selection) result = delivery.run_plans([plan], [], None) outbox = out_dir(root, "t") / delivery.OUTBOX assert result.sent == 1 assert [p.name for p in outbox.iterdir()] == ["5 ~ Song.mp3"] # nothing to do the second time again = delivery.plan_send("t", root, library, selection) assert again.copies == [] and again.deletes == [] # it arrived: the request is gone, so the file goes req.write_requests(in_dir(root, "t"), []) plan = delivery.plan_send("t", root, library, selection) delivery.run_plans([plan], [], None) assert list(outbox.iterdir()) == [] def test_missing_file_is_reported_not_crashed(self, tmp_path): library = Library() library.tracks[5] = Track(track_id=5, name="Gone", location="/nope.mp3") root = tmp_path / "root" req.write_requests(in_dir(root, "t"), [5]) plan = delivery.plan_send("t", root, library, Selection(all_library=True)) assert plan.missing == [5] and plan.copies == [] # ---- receiving ---- def _their_share(root, token, library, playlists=()): for playlist in playlists: library.playlists[playlist.persistent_id] = playlist share.publish_share(in_dir(root, token), library, Selection(all_library=True, all_playlists=True)) def _deliver(root, token, source: Path, tid: int): outbox = in_dir(root, token) / delivery.OUTBOX outbox.mkdir(parents=True, exist_ok=True) shutil.copyfile(source, outbox / delivery.outbox_name(tid, str(source))) class TestReceive: def _setup(self, tmp_path, mp3_file, **friend_kw): theirs = _library_with_file(tmp_path, mp3_file) root = tmp_path / "root" _their_share(root, "t", theirs, friend_kw.pop("playlists", ())) friend = Friend(token="t", device_id="D", name="Sam", **friend_kw) return theirs, root, friend def test_cassetted_file_is_planned_for_import(self, tmp_path, mp3_file): theirs, root, friend = self._setup(tmp_path, mp3_file, cassetted=[5]) _deliver(root, "t", Path(theirs.tracks[5].location), 5) plan = delivery.plan_receive(friend, root, matching.LibraryIndex([])) [arrival] = plan.arrivals assert arrival.to_library and not arrival.to_cache assert plan.cassetted == set() and plan.requests_after == set() def test_incomplete_is_left_alone(self, tmp_path, mp3_file): theirs, root, friend = self._setup(tmp_path, mp3_file, cassetted=[5]) _deliver(root, "t", Path(theirs.tracks[5].location), 5) plan = delivery.plan_receive(friend, root, matching.LibraryIndex([]), is_complete=lambda p: False) assert plan.arrivals == [] and plan.requests_after == {5} def test_already_have_it_is_not_imported_twice(self, tmp_path, mp3_file): theirs, root, friend = self._setup(tmp_path, mp3_file, cassetted=[5]) _deliver(root, "t", Path(theirs.tracks[5].location), 5) mine = [Track(track_id=1, name="Song", artist="Artist", album="Album", total_time=1000)] plan = delivery.plan_receive(friend, root, matching.LibraryIndex(mine)) assert plan.arrivals == [] and plan.already_have == [5] assert plan.requests_after == set() def test_prunes_what_they_stopped_offering(self, tmp_path, mp3_file): theirs, root, friend = self._setup(tmp_path, mp3_file, cassetted=[5, 77]) plan = delivery.plan_receive(friend, root, matching.LibraryIndex([])) assert plan.requests_after == {5} def test_no_library_file_yet_prunes_nothing(self, tmp_path): root = tmp_path / "root" friend = Friend(token="t", device_id="D", name="Sam", cassetted=[5, 77]) plan = delivery.plan_receive(friend, root, matching.LibraryIndex([])) assert plan.requests_after == {5, 77} def test_followed_playlist_goes_to_the_cache(self, tmp_path, mp3_file): theirs, root, friend = self._setup( tmp_path, mp3_file, followed=["P"], playlists=[Playlist(name="Mix", persistent_id="P", track_ids=[5])]) assert delivery.plan_receive(friend, root, matching.LibraryIndex([])).requests_after == {5} _deliver(root, "t", Path(theirs.tracks[5].location), 5) plan = delivery.plan_receive(friend, root, matching.LibraryIndex([])) [arrival] = plan.arrivals assert arrival.to_cache and not arrival.to_library result = delivery.run_plans([], [plan], tmp_path / "Music") assert (cache_dir(root, "t") / "5.mp3").exists() assert result.staged == [] # cached: no longer asked for again = delivery.plan_receive(friend, root, matching.LibraryIndex([])) assert again.requests_after == set() def test_followed_into_library_mode(self, tmp_path, mp3_file): theirs, root, friend = self._setup( tmp_path, mp3_file, followed=["P"], followed_mode=FOLLOW_LIBRARY, playlists=[Playlist(name="Mix", persistent_id="P", track_ids=[5])]) _deliver(root, "t", Path(theirs.tracks[5].location), 5) [arrival] = delivery.plan_receive(friend, root, matching.LibraryIndex([])).arrivals assert arrival.to_library and arrival.to_cache def test_staging_uses_the_original_name(self, tmp_path, mp3_file): theirs, root, friend = self._setup(tmp_path, mp3_file, cassetted=[5]) _deliver(root, "t", Path(theirs.tracks[5].location), 5) plan = delivery.plan_receive(friend, root, matching.LibraryIndex([])) result = delivery.run_plans([], [plan], tmp_path / "Music") [staged] = result.staged assert staged.dest.name == "Song.mp3" assert tmp_path / "Music" in staged.dest.parents # ---- the whole loop, Syncthing played by shutil ---- def _sync(a_root, b_root, token): """Copy each side's out/ to the other's in/, like Syncthing would (including deletions).""" for src_root, dst_root in ((a_root, b_root), (b_root, a_root)): src, dst = out_dir(src_root, token), in_dir(dst_root, token) if dst.exists(): shutil.rmtree(dst) if src.exists(): shutil.copytree(src, dst, ignore=shutil.ignore_patterns("*.part", "*.tmp")) def test_loop_between_two_libraries(qapp, tmp_path, mp3_file): from lintunes.gui.cassette_delivery import DeliveryCoordinator from lintunes.library_manager import LibraryManager token = "0123456789abcdef" alice_lib = _library_with_file(tmp_path, mp3_file, tid=5, name="Doo Wop", artist="Lauryn Hill", album="Miseducation") alice_mgr = LibraryManager(alice_lib, tmp_path / "alice-data") bob_mgr = LibraryManager(Library(), tmp_path / "bob-data") bob_mgr.organize_root = lambda: tmp_path / "bob-music" def service(name, manager, selection=None): s = CassetteService(root=tmp_path / name, config_loader=lambda: {}) s._state.friends[token] = Friend(token=token, device_id=name.upper(), name=name, in_folder="x", selection=selection or Selection()) s.set_library_provider(lambda: manager.library) s.is_complete = lambda t, p: True return s alice = service("alice", alice_mgr, Selection(all_library=True)) bob = service("bob", bob_mgr) alice.republish() _sync(alice.root, bob.root, token) # Bob cassettes Alice's song and saves. bob.save_requests(token, {5}, []) _sync(alice.root, bob.root, token) def run(coordinator): coordinator.run() deadline = time.time() + 10 while coordinator.busy() and time.time() < deadline: qapp.processEvents() time.sleep(0.01) qapp.processEvents() a_loop = DeliveryCoordinator(alice, alice_mgr) b_loop = DeliveryCoordinator(bob, bob_mgr) arrived = [] b_loop.notice.connect(arrived.append) run(a_loop) # Alice sends assert [p.name for p in (out_dir(alice.root, token) / "outbox").iterdir()] \ == ["5 ~ Doo Wop.mp3"] _sync(alice.root, bob.root, token) run(b_loop) # Bob brings it in [track] = bob_mgr.library.tracks.values() assert track.name != "" and Path(track.location).exists() assert tmp_path / "bob-music" in Path(track.location).parents assert track.play_count == 0 and track.rating == 0 assert arrived == ["1 song arrived from bob."] assert req.read_requests(out_dir(bob.root, token)) == set() assert bob.state().friends[token].cassetted == [] _sync(alice.root, bob.root, token) run(a_loop) # Alice's outbox cleans up assert list((out_dir(alice.root, token) / "outbox").iterdir()) == [] _sync(alice.root, bob.root, token) run(b_loop) # and nothing is imported twice assert len(bob_mgr.library.tracks) == 1 # ---- the real thing ---- @pytest.mark.syncthing def test_real_delivery(qapp, tmp_path, syncthing_pair, mp3_file): from lintunes.cassette import invite as inv from lintunes.cassette import service as svc_mod from lintunes.cassette import syncthing_api as st from lintunes.gui.cassette_delivery import DeliveryCoordinator from lintunes.library_manager import LibraryManager alice_i, bob_i, _carol = syncthing_pair listen = {i.device_id: [i.listen_address] for i in syncthing_pair} alice_mgr = LibraryManager( _library_with_file(tmp_path, mp3_file, tid=5, name="Doo Wop", artist="Lauryn Hill", album="Miseducation"), tmp_path / "alice-data") bob_mgr = LibraryManager(Library(), tmp_path / "bob-data") bob_mgr.organize_root = lambda: tmp_path / "bob-music" def make(instance, manager): config = {st.API_KEY_KEY: instance.api_key, st.ADDRESS_KEY: instance.address} service = CassetteService(root=tmp_path / instance.name, config_loader=lambda: config, addresses_for=lambda d: listen[d]) service._state.display_name = instance.name.title() service.set_library_provider(lambda: manager.library) return service, service.client(), DeliveryCoordinator(service, manager) a, a_client, a_loop = make(alice_i, alice_mgr) b, b_client, b_loop = make(bob_i, bob_mgr) code = a.create_invite() b.accept_code(code) token = inv.decode(code).token def settle(condition, timeout=120): deadline = time.time() + timeout while time.time() < deadline: for service, client in ((a, a_client), (b, b_client)): service.reconcile(client) for loop in (a_loop, b_loop): if not loop.busy(): loop.run() qapp.processEvents() if condition(): return True time.sleep(0.5) return False assert settle(lambda: token in a.state().friends and a.status(token) and a.status(token).code == svc_mod.CONNECTED) edited = a.state() edited.friends[token].selection = Selection(all_library=True) a.apply_settings(edited) assert settle(lambda: 5 in load_their_ids(b, token)) b.save_requests(token, {5}, []) assert settle(lambda: len(bob_mgr.library.tracks) == 1), "never arrived" assert settle(lambda: not any((out_dir(a.root, token) / "outbox").glob("*"))), \ "sender's outbox never cleaned up" def load_their_ids(service, token): from lintunes.cassette.friend_library import load_friend_library return set(load_friend_library(in_dir(service.root, token)).tracks)