v0.30.0: Cassette delivery — requested songs arrive in your library
- Sending: a friend's requests ∩ what I still offer is copied into my outbox (via a .part name Syncthing never sees); whatever they no longer request is deleted, so the outbox cleans itself up. - Receiving: a finished file I cassetted is imported exactly like Add to Library — organized into Artist/Album, tags read, date added now, none of their listening history — and a song wanted only for a followed playlist lands in the cache. Then requests.json is rewritten without what arrived and without what they stopped offering. - Safe to repeat: Syncthing temp files ignored, Syncthing asked whether a file is whole, and an exact match already in my library is never imported twice. - Folders are watched (re-armed after each rename), with a slow fallback. - file_importer split into stage_file (off the GUI thread) and track_from_file (date added = now). - pytest --syncthing now runs a real invite → share → request → deliver → cleanup between two Syncthing instances. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
@@ -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/<id> ~ <name>` (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
|
||||
|
||||
@@ -1,3 +1,3 @@
|
||||
"""LinTunes — iTunes-style music library manager and player for Linux."""
|
||||
|
||||
__version__ = "0.29.0"
|
||||
__version__ = "0.30.0"
|
||||
|
||||
@@ -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
|
||||
@@ -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:
|
||||
|
||||
@@ -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)
|
||||
@@ -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()
|
||||
|
||||
@@ -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()
|
||||
|
||||
@@ -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 ``<id> ~ <name>``)."""
|
||||
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
|
||||
|
||||
|
||||
|
||||
@@ -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)
|
||||
Reference in New Issue
Block a user