diff --git a/CLAUDE.md b/CLAUDE.md index dc9fe78..18ac0c4 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -493,6 +493,18 @@ persistence) → GUI (Qt widgets that read the manager and connect to its signal asks "Save changes?". `cassette/matching.py` is best effort (exact vs fuzzy); `cassette/space.py` judges the music disk (red) and Syncthing's `minDiskFree` floor on the cassette disk (orange) separately. + **Delivery** (`cassette/delivery.py` + `gui/cassette_delivery.py`) is + rebuilt from the folders on every pass, both directions: sender = their + `requests.json` ∩ what I offer *now* → `outbox/ ~ ` (copied to a + `.part` the `.stignore` hides, then renamed), and anything no longer + requested is deleted — a request vanishing *is* the delivery receipt. + Receiver: planning on the GUI thread (it reads the library), copying on + `DeliveryWorker`'s thread (`file_importer.stage_file` — the import split so + the copy is off the GUI thread), `manager.add_track` back on the GUI thread + (`track_from_file`: date added now, no listening history). An exact match + already in my library is never imported twice, which is what makes a crash + between import and requests-rewrite harmless; a missing `library.json` + prunes nothing. - **`lintunes/cast/`** — Chromecast playback (Connections menu), using the **media-receiver model**: `server.py` runs a `ThreadingHTTPServer` on an diff --git a/lintunes/__init__.py b/lintunes/__init__.py index 1a0b003..264d54c 100644 --- a/lintunes/__init__.py +++ b/lintunes/__init__.py @@ -1,3 +1,3 @@ """LinTunes — iTunes-style music library manager and player for Linux.""" -__version__ = "0.29.0" +__version__ = "0.30.0" diff --git a/lintunes/cassette/delivery.py b/lintunes/cassette/delivery.py new file mode 100644 index 0000000..43547fc --- /dev/null +++ b/lintunes/cassette/delivery.py @@ -0,0 +1,305 @@ +"""The request / delivery loop — both directions, rebuilt from the folders. + +**Sending** (my LinTunes serving a friend): their ``requests.json`` (in their +folder to me) ∩ what I *currently* offer them = what belongs in my +``outbox/``. Missing files are copied in (to a ``.part`` name the +``.stignore`` keeps out of sync, then renamed); files no longer requested are +deleted — a request disappears when the song has arrived, so the outbox +cleans itself up with no bookkeeping of its own. + +**Receiving** (my LinTunes asking a friend): each complete file in their +outbox that I want is brought in — a cassetted song is imported exactly like +Add to Library (copied into the organized tree, tags read, date added now, +none of their listening history), a song wanted only for a followed playlist +is copied into ``cache/``. Then my ``requests.json`` is rewritten once per +batch without what arrived and without anything they no longer offer. + +Each machine trusts only its own view, and every step is safe to repeat: a +song already in my library (exact match) is not imported twice, a file +Syncthing is still writing (``.syncthing.*.tmp``) is never touched, and +whatever a crash interrupted is simply planned again next time. + +Planning is pure and runs on the GUI thread (it reads the library); the file +work runs on ``DeliveryWorker``'s thread; adding tracks to the library happens +back on the GUI thread. +""" +from __future__ import annotations + +import logging +import re +import shutil +import threading +from dataclasses import dataclass, field +from pathlib import Path + +from PyQt6.QtCore import QObject, pyqtSignal + +from lintunes.cassette import matching, publish, share +from lintunes.cassette import requests as req +from lintunes.cassette.friend_library import load_friend_library +from lintunes.cassette.state import FOLLOW_LIBRARY, cache_dir, in_dir, out_dir +from lintunes.importers.file_importer import sanitize_component + +log = logging.getLogger(__name__) + +OUTBOX = "outbox" +PART_SUFFIX = ".part" +_NAME_RE = re.compile(r"^(\d+) ~ (.+)$") + + +def outbox_name(track_id: int, location: str) -> str: + return f"{track_id} ~ {sanitize_component(Path(location).name)}" + + +def parse_outbox_name(name: str) -> tuple[int, str] | None: + """(track id, original filename) for one of ours, else None.""" + match = _NAME_RE.match(name) + return (int(match.group(1)), match.group(2)) if match else None + + +def is_syncthing_temp(name: str) -> bool: + return (name.startswith(".syncthing.") or name.startswith("~syncthing~") + or name.endswith(".tmp") or name.endswith(PART_SUFFIX)) + + +def list_outbox(folder: Path) -> dict[int, Path]: + """{track id: path} for the finished files in an outbox.""" + found = {} + try: + entries = list(folder.iterdir()) + except OSError: + return found + for path in entries: + if is_syncthing_temp(path.name) or not path.is_file(): + continue + parsed = parse_outbox_name(path.name) + if parsed is not None: + found[parsed[0]] = path + return found + + +# ---------------------------------------------------------------- sending -- + +@dataclass +class SendPlan: + token: str + outbox: Path + copies: list = field(default_factory=list) # (source Path, dest name) + deletes: list = field(default_factory=list) # Paths + missing: list = field(default_factory=list) # track ids with no file + + +def plan_send(token: str, root: Path, library, selection) -> SendPlan: + outbox = out_dir(root, token) / OUTBOX + their_requests = req.read_requests(in_dir(root, token)) + offered = share.offered_ids(library, selection) + want = their_requests & offered + present = list_outbox(outbox) + plan = SendPlan(token=token, outbox=outbox) + for tid in sorted(want): + if tid in present: + continue + track = library.tracks[tid] + source = Path(track.location) if track.location else None + if source is None or not source.is_file(): + plan.missing.append(tid) + continue + plan.copies.append((source, outbox_name(tid, track.location))) + for tid, path in present.items(): + if tid not in want: + plan.deletes.append(path) + return plan + + +# -------------------------------------------------------------- receiving -- + +@dataclass +class Arrival: + track_id: int + path: Path + to_library: bool + to_cache: bool + + +@dataclass +class ReceivePlan: + token: str + arrivals: list = field(default_factory=list) # Arrival + cache: Path = None + # The requests file as it should read once the arrivals are in. + cassetted: set = field(default_factory=set) + requests_after: set = field(default_factory=set) + requests_path: Path = None + already_have: list = field(default_factory=list) # ids skipped: exact match + + +def followed_needed(friend_library, friend, cache: Path) -> set[int]: + """Songs of followed playlists that aren't in the cache yet.""" + cached = {int(p.stem) for p in _cache_files(cache)} + needed = set() + for pid in friend.followed: + playlist = friend_library.playlist(pid) + if playlist is not None: + needed.update(t for t in playlist.track_ids if t not in cached) + return needed + + +def _cache_files(cache: Path): + try: + return [p for p in cache.iterdir() + if p.is_file() and p.stem.isdigit() and not is_syncthing_temp(p.name)] + except OSError: + return [] + + +def plan_receive(friend, root: Path, my_index: matching.LibraryIndex, + is_complete=None) -> ReceivePlan: + token = friend.token + source = in_dir(root, token) + library_file = source / share.LIBRARY_FILE + their = load_friend_library(source, friend.name) + cache = cache_dir(root, token) + plan = ReceivePlan(token=token, cache=cache, + requests_path=out_dir(root, token)) + + cassetted = set(friend.cassetted) + # Prune what they stopped offering — but only against a library file + # that's actually here; a not-yet-synced one isn't "offers nothing". + if library_file.exists() and publish.read_json(library_file) is not None: + cassetted &= set(their.tracks) + to_library = set(cassetted) + needed = followed_needed(their, friend, cache) + if friend.followed_mode == FOLLOW_LIBRARY: + to_library |= needed + needed_cache = set(needed) + else: + needed_cache = set(needed) + + for tid, path in sorted(list_outbox(source / OUTBOX).items()): + wanted_lib = tid in to_library + wanted_cache = tid in needed_cache + if not (wanted_lib or wanted_cache): + continue + if is_complete is not None and not is_complete(path): + continue + if wanted_lib and tid in their.tracks and \ + my_index.classify(their.tracks[tid]) == matching.EXACT: + plan.already_have.append(tid) + wanted_lib = False + if not wanted_cache: + cassetted.discard(tid) + continue + plan.arrivals.append(Arrival(tid, path, wanted_lib, wanted_cache)) + + arrived_lib = {a.track_id for a in plan.arrivals if a.to_library} + arrived_cache = {a.track_id for a in plan.arrivals if a.to_cache} + plan.cassetted = cassetted - arrived_lib + still_needed = needed - arrived_cache + if friend.followed_mode == FOLLOW_LIBRARY: + still_needed -= arrived_lib + plan.requests_after = plan.cassetted | still_needed + return plan + + +# ------------------------------------------------------------------ work -- + +@dataclass +class Staged: + token: str + track_id: int + dest: Path + fields: dict + fallback_name: str + + +@dataclass +class DeliveryResult: + staged: list = field(default_factory=list) # Staged, to add to library + cached: list = field(default_factory=list) # (token, track id) + sent: int = 0 + removed: int = 0 + errors: list = field(default_factory=list) # (what, error text) + + +def run_plans(send_plans, receive_plans, music_dir: Path | None, + cancel: threading.Event | None = None) -> DeliveryResult: + """All the file work. Safe off the GUI thread: touches no library.""" + from lintunes.importers.file_importer import stage_file + result = DeliveryResult() + for plan in send_plans: + for path in plan.deletes: + try: + path.unlink(missing_ok=True) + result.removed += 1 + except OSError as exc: + result.errors.append((f"Removing {path.name}", str(exc))) + if plan.copies: + plan.outbox.mkdir(parents=True, exist_ok=True) + for source, name in plan.copies: + if cancel is not None and cancel.is_set(): + return result + part = plan.outbox / ("." + name + PART_SUFFIX) + try: + shutil.copyfile(source, part) + part.replace(plan.outbox / name) + result.sent += 1 + except OSError as exc: + part.unlink(missing_ok=True) + result.errors.append((f"Sending {source.name}", str(exc))) + for plan in receive_plans: + for arrival in plan.arrivals: + if cancel is not None and cancel.is_set(): + return result + parsed = parse_outbox_name(arrival.path.name) + original = parsed[1] if parsed else arrival.path.name + try: + if arrival.to_cache: + plan.cache.mkdir(parents=True, exist_ok=True) + dest = plan.cache / f"{arrival.track_id}{arrival.path.suffix.lower()}" + tmp = dest.with_name("." + dest.name + PART_SUFFIX) + shutil.copyfile(arrival.path, tmp) + tmp.replace(dest) + result.cached.append((plan.token, arrival.track_id)) + if arrival.to_library and music_dir is not None: + dest, fields = stage_file(arrival.path, music_dir, original) + result.staged.append(Staged(plan.token, arrival.track_id, + dest, fields, + Path(original).stem)) + except OSError as exc: + result.errors.append((f"Bringing in {original}", str(exc))) + return result + + +class DeliveryWorker(QObject): + """``run_plans`` on a daemon thread (the ExportWorker shape).""" + + finished = pyqtSignal(object, object) # DeliveryResult, receive plans + + def __init__(self, parent=None): + super().__init__(parent) + self._busy = False + self._cancel = threading.Event() + + def busy(self) -> bool: + return self._busy + + def cancel(self): + self._cancel.set() + + def start(self, send_plans, receive_plans, music_dir): + if self._busy: + return False + self._busy = True + self._cancel.clear() + + def run(): + try: + result = run_plans(send_plans, receive_plans, music_dir, + self._cancel) + except Exception as exc: # never die silently + log.exception("cassette delivery") + result = DeliveryResult(errors=[("Delivering songs", str(exc))]) + self._busy = False + self.finished.emit(result, receive_plans) + threading.Thread(target=run, daemon=True, name="cassette-delivery").start() + return True diff --git a/lintunes/cassette/service.py b/lintunes/cassette/service.py index 15ff826..1b7f3ea 100644 --- a/lintunes/cassette/service.py +++ b/lintunes/cassette/service.py @@ -254,6 +254,42 @@ class CassetteService(QObject): req.write_requests(out_dir(self.root, token), requested) self.changed.emit() + def apply_delivery(self, token: str, cassetted, requested): + """After a batch arrived: what's still wanted, and the requests file + that says so (rewritten only if it changed).""" + from lintunes.cassette import requests as req + with self._lock: + friend = self._state.friends.get(token) + if friend is None: + return + cassetted = sorted(int(t) for t in cassetted) + if cassetted != friend.cassetted: + friend.cassetted = cassetted + self._save() + req.write_requests(out_dir(self.root, token), requested) + + def is_complete(self, token: str, path) -> bool: + """Has Syncthing finished writing this file of the friend's? It only + renames a download into place once it's whole, so a finished name is + normally enough; this double-checks with Syncthing when it can.""" + with self._lock: + friend = self._state.friends.get(token) + if friend is None: + return False + try: + client = self.client() + relative = path.relative_to(in_dir(self.root, token)).as_posix() + info = client.file_info(friend.in_folder, relative) + except (st.SyncthingError, ValueError): + return True + if not info: + return True + local, global_ = info.get("local") or {}, info.get("global") or {} + if local.get("deleted") or global_.get("deleted"): + return False + return (local.get("version") == global_.get("version") + and local.get("size") == global_.get("size")) + # ---- user actions ---- def create_invite(self) -> str: diff --git a/lintunes/gui/cassette_delivery.py b/lintunes/gui/cassette_delivery.py new file mode 100644 index 0000000..00e33b3 --- /dev/null +++ b/lintunes/gui/cassette_delivery.py @@ -0,0 +1,143 @@ +"""Runs the delivery loop (``cassette/delivery.py``) inside the app. + +Watches every friend's folder (their requests, their outbox) with a +QFileSystemWatcher — re-armed after each event, since Syncthing's +rename-into-place drops the watch, the ``sync_watcher`` lesson — and a slow +fallback timer for whatever a watcher misses. Each trigger is debounced into +one pass: plan on the GUI thread, copy on the worker's thread, add tracks to +the library back on the GUI thread, then rewrite requests. +""" +from __future__ import annotations + +from pathlib import Path + +from PyQt6.QtCore import QFileSystemWatcher, QObject, QTimer, pyqtSignal + +from lintunes.cassette import delivery, matching +from lintunes.cassette.state import in_dir + +DEBOUNCE_MS = 2000 +FALLBACK_MS = 10 * 60 * 1000 + + +class DeliveryCoordinator(QObject): + notice = pyqtSignal(str) + step_failed = pyqtSignal(str, str) + library_grew = pyqtSignal(list) # new track ids + + def __init__(self, service, manager, parent=None): + super().__init__(parent) + self._service = service + self._manager = manager + self._worker = delivery.DeliveryWorker(self) + self._worker.finished.connect(self._on_finished) + self._again = False + self._watcher = QFileSystemWatcher(self) + self._watcher.directoryChanged.connect(self._on_change) + self._watcher.fileChanged.connect(self._on_change) + self._debounce = QTimer(self) + self._debounce.setSingleShot(True) + self._debounce.setInterval(DEBOUNCE_MS) + self._debounce.timeout.connect(self.run) + self._fallback = QTimer(self) + self._fallback.setInterval(FALLBACK_MS) + self._fallback.timeout.connect(self.run) + + def start(self): + self._arm() + self._fallback.start() + self._debounce.start() + + def stop(self): + self._fallback.stop() + self._debounce.stop() + self._worker.cancel() + + def busy(self) -> bool: + return self._worker.busy() + + def _arm(self): + """(Re)watch every friend's folder and outbox that exists.""" + watched = set(self._watcher.directories()) | set(self._watcher.files()) + wanted = [] + for token in self._service.state().friends: + folder = in_dir(self._service.root, token) + for path in (folder, folder / delivery.OUTBOX, + folder / "requests.json"): + if path.exists() and str(path) not in watched: + wanted.append(str(path)) + if wanted: + self._watcher.addPaths(wanted) + + def _on_change(self, _path: str): + self._arm() + self._debounce.start() + + def schedule(self): + self._debounce.start() + + def run(self): + """Plan every friend in both directions and start the file work.""" + if self._worker.busy(): + self._again = True + return + self._arm() + state = self._service.state() + if not state.friends: + return + library = self._manager.library + index = matching.LibraryIndex(library.tracks.values()) + root = self._service.root + sends, receives = [], [] + for token, friend in state.friends.items(): + sends.append(delivery.plan_send(token, root, library, + state.selection_for(friend))) + receives.append(delivery.plan_receive( + friend, root, index, + is_complete=lambda path, t=token: self._service.is_complete(t, path))) + if not any(p.copies or p.deletes for p in sends) and \ + not any(p.arrivals for p in receives): + # Nothing to move, but prunes and exact-match skips still count. + self._finish_receives(receives, set()) + return + self._worker.start(sends, receives, self._manager.organize_root()) + + def _on_finished(self, result, receives): + added = [] + for staged in result.staged: + from lintunes.importers.file_importer import track_from_file + track = track_from_file(staged.dest, staged.fields, + staged.fallback_name) + self._manager.add_track(track) + added.append(track.track_id) + brought = {(s.token, s.track_id) for s in result.staged} + brought |= set(result.cached) + self._finish_receives(receives, brought) + if added: + self.library_grew.emit(added) + names = {t: f.name for t, f in self._service.state().friends.items()} + per_friend = {} + for token, _tid in brought: + per_friend[token] = per_friend.get(token, 0) + 1 + for token, count in per_friend.items(): + noun = "song" if count == 1 else "songs" + self.notice.emit(f"{count} {noun} arrived from " + f"{names.get(token, 'a friend')}.") + for what, error in result.errors[:3]: + self.step_failed.emit(what, error) + if self._again: + self._again = False + self._debounce.start() + + def _finish_receives(self, receives, brought: set): + """Rewrite each friend's requests without what arrived (only an + arrival that really landed counts — a failed copy is asked for + again).""" + for plan in receives: + failed = {a.track_id for a in plan.arrivals + if (plan.token, a.track_id) not in brought} + self._service.apply_delivery( + plan.token, plan.cassetted | {t for t in failed + if any(a.track_id == t and a.to_library + for a in plan.arrivals)}, + plan.requests_after | failed) diff --git a/lintunes/gui/cassette_ui.py b/lintunes/gui/cassette_ui.py index 4245dd9..300f774 100644 --- a/lintunes/gui/cassette_ui.py +++ b/lintunes/gui/cassette_ui.py @@ -37,6 +37,7 @@ class CassetteUi(QObject): self._service_factory = service_factory self.service = None self.friend_mode = None # FriendMode while browsing a friend + self.delivery = None # DeliveryCoordinator (host machine only) # ---- the service (host machine only) ---- @@ -70,6 +71,15 @@ class CassetteUi(QObject): self.service.start() # Whatever changed while LinTunes was closed. self._republish_timer.start() + if self._manager is not None: + from lintunes.gui.cassette_delivery import DeliveryCoordinator + self.delivery = DeliveryCoordinator(self.service, self._manager, self) + self.delivery.notice.connect(self._on_notice) + self.delivery.step_failed.connect(self._on_step_failed) + self.delivery.library_grew.connect(self._on_library_grew) + self.service.friends_added.connect(self.delivery.schedule) + self.delivery.start() + self._refresh_selector() return self.service def _library(self): @@ -134,6 +144,9 @@ class CassetteUi(QObject): self.friend_mode = FriendMode( self._window, self.service, self._manager.library, music_dir or self.service.root) + if self.delivery is not None: + # A save changed what I ask for: act on it now, not in ten minutes. + self.friend_mode.left.connect(self.delivery.schedule) self.friend_mode.enter(token) def leave_friend_mode(self): @@ -145,7 +158,14 @@ class CassetteUi(QObject): return True return self.friend_mode.confirm_leave() + def _on_library_grew(self, track_ids): + refresh = getattr(self._window, "refresh_library_view", None) + if refresh is not None: + refresh(track_ids) + def shutdown(self): + if self.delivery is not None: + self.delivery.stop() if self.service is not None: self.service.stop() @@ -263,3 +283,5 @@ class CassetteUi(QObject): except st.SyncthingError as exc: QMessageBox.warning(self._window, "Sync Settings", failure_text(exc)) service.kick() + if self.delivery is not None: + self.delivery.schedule() diff --git a/lintunes/gui/main_window.py b/lintunes/gui/main_window.py index c309dee..6eb6bf6 100644 --- a/lintunes/gui/main_window.py +++ b/lintunes/gui/main_window.py @@ -403,6 +403,11 @@ class MainWindow(QMainWindow): def sidebar(self): return self._sidebar + def refresh_library_view(self, _track_ids=None): + """Songs arrived from somewhere other than an import dialog (a friend's + delivery): show them.""" + self._library_view.reload() + def show_my_library(self): self._sidebar.reflect_library() self._show_library() diff --git a/lintunes/importers/file_importer.py b/lintunes/importers/file_importer.py index c3cb928..f2299eb 100644 --- a/lintunes/importers/file_importer.py +++ b/lintunes/importers/file_importer.py @@ -79,9 +79,36 @@ def import_file(source: Path, music_dir: Path, manager, dest.parent.mkdir(parents=True, exist_ok=True) shutil.copy2(source, dest) + track = track_from_file(dest, fields, source.stem) + manager.add_track(track) + by_location[Path(track.location)] = track + return track + + +def stage_file(source: Path, music_dir: Path, + filename: str | None = None) -> tuple[Path, dict]: + """The file-system half of an import, safe off the GUI thread: read the + tags and copy the file to its organized home. Returns (dest, fields); + ``track_from_file`` + ``manager.add_track`` finish it on the GUI thread. + ``filename`` overrides the name the copy gets (a Cassette delivery + arrives as `` ~ ``).""" + source = Path(source) + fields = tagging.read_tags(source) + dest = unique_path(organized_destination(music_dir, fields, + filename or source.name)) + dest.parent.mkdir(parents=True, exist_ok=True) + shutil.copy2(source, dest) + return dest, fields + + +def track_from_file(dest: Path, fields: dict, fallback_name: str) -> Track: + """A new Track for a file just imported. Date added is *now* — every + individual import is new to this library, whatever the file carries; + only the iTunes migration keeps old dates, and it doesn't come here. + Only the file's own tags come across: nobody's listening history.""" now = datetime.now(timezone.utc).replace(tzinfo=None).isoformat() track = Track( - name=fields.get("name") or source.stem, + name=fields.get("name") or fallback_name, location=str(dest), date_added=now, date_modified=now, @@ -91,8 +118,6 @@ def import_file(source: Path, music_dir: Path, manager, "comments", "total_time", "bit_rate", "sample_rate", "size", "kind"): if key in fields: setattr(track, key, fields[key]) - manager.add_track(track) - by_location[Path(track.location)] = track return track diff --git a/tests/test_round68.py b/tests/test_round68.py new file mode 100644 index 0000000..82c045f --- /dev/null +++ b/tests/test_round68.py @@ -0,0 +1,337 @@ +"""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}, [], {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}, [], {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)