10 Commits
33 changed files with 2884 additions and 141 deletions
+40
View File
@@ -0,0 +1,40 @@
# Changelog
All notable changes to Furumi are documented in this file.
The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.1.0/),
and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0.html).
## [Unreleased]
## [0.2.5] - 2026-08-02
### Added
- Command history navigation with the Up and Down arrow keys.
- A copy-friendly plain terminal view for connection tickets, device and Jam
invites, and track-sharing links.
- A configurable permanent music directory with write validation and an
optional safe migration of existing Furumi-managed files.
- Queue reordering for one track or a `Shift+V` selection with `Alt+K` and
`Alt+J`.
- Reproducible Nix/devenv tooling with Rust 1.97 and ALSA development files.
### Changed
- Sequential `a` actions now build one ordered play-next block instead of
reversing independently inserted tracks.
- Federated artist images and release covers are stored beside permanent music
in an `Artist/Release` directory tree.
### Fixed
- Running a development build next to another Furumi instance no longer lets
an MPRIS name collision disable terminal raw mode and freeze keyboard input.
- `nix develop` now keeps devenv state outside the read-only Nix store and
isolates Rust 1.97 build artifacts from other toolchains.
- Music-directory validation and migration now reject overlapping changes,
resolve canonical paths, and produce Windows-portable managed filenames.
[Unreleased]: https://gt.hexor.cy/ab/furumi_tui/compare/v0.2.5...HEAD
[0.2.5]: https://gt.hexor.cy/ab/furumi_tui/compare/v0.2.4...v0.2.5
Generated
+12 -12
View File
@@ -434,8 +434,6 @@ dependencies = [
[[package]]
name = "block"
version = "0.1.6"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "0d8c1fef690941d3e7788d328517591fecc684c084084702d6ff1641e993699a"
[[package]]
name = "block-buffer"
@@ -1211,13 +1209,13 @@ dependencies = [
[[package]]
name = "displaydoc"
version = "0.2.6"
version = "0.2.7"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "1ac70aa55017e108007fbaf5aa0f54b021c98f92ff8af59d42eda9da96e3dd4f"
checksum = "c6232dd377dcc64799954cbd3a9bb882e9cdc1308ccd87b1c098f1fb2eaf82a8"
dependencies = [
"proc-macro2",
"quote",
"syn 2.0.119",
"syn 3.0.3",
]
[[package]]
@@ -1454,8 +1452,9 @@ dependencies = [
[[package]]
name = "federation-net"
version = "0.1.0"
source = "git+https://gt.hexor.cy/ab/frid.git?rev=8de7d1292708fa0b225e5a4a9d5ab4f0676202d3#8de7d1292708fa0b225e5a4a9d5ab4f0676202d3"
version = "0.2.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "15a8707baeccb46b5935138f9cb3df3c988c0730b807a2634d901f26b39250d6"
dependencies = [
"blake3",
"data-encoding",
@@ -1566,7 +1565,7 @@ dependencies = [
[[package]]
name = "furumi_tui"
version = "0.2.2"
version = "0.2.5"
dependencies = [
"anyhow",
"blake3",
@@ -2988,8 +2987,9 @@ dependencies = [
[[package]]
name = "music-dht"
version = "0.2.0"
source = "git+https://gt.hexor.cy/ab/frid.git?rev=8de7d1292708fa0b225e5a4a9d5ab4f0676202d3#8de7d1292708fa0b225e5a4a9d5ab4f0676202d3"
version = "0.3.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "5f32daa9edf769fb5686e92ae6884e9fda6ea082452e4208309fc14301e26aef"
dependencies = [
"async-trait",
"blake3",
@@ -5578,9 +5578,9 @@ dependencies = [
[[package]]
name = "toml"
version = "1.1.3+spec-1.1.0"
version = "1.1.4+spec-1.1.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "53c96ecdfa941c8fc4fcaed14f99ada8ebed502eef533015095a07e3301d4c3c"
checksum = "3aace63f4bbcdfc2c965b059de67119c89c4017a70d633be6c104910f67056f5"
dependencies = [
"indexmap",
"serde_core",
+7 -2
View File
@@ -1,6 +1,6 @@
[package]
name = "furumi_tui"
version = "0.2.2"
version = "0.2.5"
edition = "2024"
rust-version = "1.97"
description = "A federated P2P player for personal music libraries"
@@ -21,7 +21,7 @@ image = { version = "0.25.10", default-features = false, features = ["jpeg", "pn
lofty = "0.22"
# P2P federation: library index in a shared DHT + audio streaming between
# peers (same protocol as furumi-fd).
music-dht = { git = "https://gt.hexor.cy/ab/frid.git", rev = "8de7d1292708fa0b225e5a4a9d5ab4f0676202d3" }
music-dht = "0.3"
ratatui = "0.30.1"
rhai = { version = "1", features = ["sync"] }
rodio = { version = "0.22.2", default-features = false, features = ["playback", "mp3", "flac", "vorbis", "wav", "symphonia-aac", "symphonia-isomp4", "symphonia-alac"] }
@@ -36,6 +36,11 @@ tracing = "0.1.44"
tracing-subscriber = { version = "0.3.23", features = ["env-filter"] }
unicode-width = "0.2.2"
[patch.crates-io]
# souvlaki 0.8.3 still depends on the unmaintained block 0.1.6. Keep its API
# intact while using an opaque inhabited FFI type accepted by current Rust.
block = { path = "vendor/block" }
[target.'cfg(target_os="macos")'.dependencies]
core-foundation = "0.10.1"
+40
View File
@@ -0,0 +1,40 @@
{ config, lib, pkgs, ... }:
let
sourceRoot = builtins.toString ./.;
sourceKey = builtins.substring 0 12 (builtins.hashString "sha256" sourceRoot);
transientDotfile = "/tmp/furumi-tui-devenv-${sourceKey}";
rustToolchain = builtins.fromTOML (builtins.readFile ./rust-toolchain.toml);
rustVersion = rustToolchain.toolchain.channel;
in
{
# Pure flakes only expose their read-only store source. Keep devenv's own
# task/profile state writable, then restore the real checkout path in-shell.
devenv.root = sourceRoot;
devenv.dotfile = transientDotfile;
devenv.state = "${transientDotfile}/state";
languages.rust = {
enable = true;
toolchainFile = ./rust-toolchain.toml;
};
packages = [
pkgs.pkg-config
] ++ pkgs.lib.optionals pkgs.stdenv.isLinux [
pkgs.alsa-lib
];
env.RUST_BACKTRACE = "1";
enterShell = lib.mkAfter ''
export DEVENV_ROOT="$PWD"
# Keep Nix builds separate from artifacts produced by rustup or another
# devenv generation. Cargo metadata is not compatible across compilers.
export CARGO_TARGET_DIR="$PWD/target/devenv-rust-${rustVersion}"
export PATH="${config.languages.rust.toolchainPackage}/bin:$PATH"
export RUSTC="${config.languages.rust.toolchainPackage}/bin/rustc"
export RUSTDOC="${config.languages.rust.toolchainPackage}/bin/rustdoc"
hash -r
'';
}
Generated
+399
View File
@@ -0,0 +1,399 @@
{
"nodes": {
"cachix": {
"inputs": {
"devenv": [
"devenv"
],
"flake-compat": [
"devenv",
"flake-compat"
],
"git-hooks": [
"devenv",
"git-hooks"
],
"nixpkgs": "nixpkgs"
},
"locked": {
"lastModified": 1777487137,
"narHash": "sha256-TuvKVBX60mqyMT6OB5JqVEh1YIWtFMR/igLCaCdC9tw=",
"owner": "cachix",
"repo": "cachix",
"rev": "a66a440c321d35f7193472c317f42a55ccd1cb93",
"type": "github"
},
"original": {
"owner": "cachix",
"ref": "latest",
"repo": "cachix",
"type": "github"
}
},
"crate2nix": {
"flake": false,
"locked": {
"lastModified": 1772186516,
"narHash": "sha256-8s28pzmQ6TOIUzznwFibtW1CMieMUl1rYJIxoQYor58=",
"owner": "rossng",
"repo": "crate2nix",
"rev": "ba5dd398e31ee422fbe021767eb83b0650303a6e",
"type": "github"
},
"original": {
"owner": "rossng",
"repo": "crate2nix",
"rev": "ba5dd398e31ee422fbe021767eb83b0650303a6e",
"type": "github"
}
},
"devenv": {
"inputs": {
"cachix": "cachix",
"crate2nix": "crate2nix",
"flake-compat": "flake-compat",
"flake-parts": "flake-parts",
"ghostty": "ghostty",
"git-hooks": "git-hooks",
"nix": "nix",
"nixd": "nixd",
"nixpkgs": [
"nixpkgs"
],
"rust-overlay": "rust-overlay"
},
"locked": {
"lastModified": 1785616853,
"narHash": "sha256-WKA53CfbSwXh/QuGHc3gsWoLNHYaKw4xgvkc9ceRgUA=",
"owner": "cachix",
"repo": "devenv",
"rev": "a3feed9ca96c66548496e8af2804d2a3a045d22e",
"type": "github"
},
"original": {
"owner": "cachix",
"repo": "devenv",
"type": "github"
}
},
"flake-compat": {
"flake": false,
"locked": {
"lastModified": 1767039857,
"narHash": "sha256-vNpUSpF5Nuw8xvDLj2KCwwksIbjua2LZCqhV1LNRDns=",
"owner": "edolstra",
"repo": "flake-compat",
"rev": "5edf11c44bc78a0d334f6334cdaf7d60d732daab",
"type": "github"
},
"original": {
"owner": "edolstra",
"repo": "flake-compat",
"type": "github"
}
},
"flake-parts": {
"inputs": {
"nixpkgs-lib": [
"devenv",
"nixpkgs"
]
},
"locked": {
"lastModified": 1778716662,
"narHash": "sha256-m1Yf0wZ8j1OHjTc2UwHwyQRSnNeSgLJOd7q5Y45hzi4=",
"owner": "hercules-ci",
"repo": "flake-parts",
"rev": "f7c1a2d347e4c52d5fb8d10cb4d94b5884e546fb",
"type": "github"
},
"original": {
"owner": "hercules-ci",
"repo": "flake-parts",
"type": "github"
}
},
"flake-parts_2": {
"inputs": {
"nixpkgs-lib": "nixpkgs-lib"
},
"locked": {
"lastModified": 1782949081,
"narHash": "sha256-vp6Y/Grm98ESt6ceOkWiHWyZRDV3J1RID4w+6NWK9yA=",
"owner": "hercules-ci",
"repo": "flake-parts",
"rev": "17c9d6cdfc60c64f4ee8d306f9bc0b4ccb51481e",
"type": "github"
},
"original": {
"owner": "hercules-ci",
"repo": "flake-parts",
"type": "github"
}
},
"ghostty": {
"flake": false,
"locked": {
"lastModified": 1784602798,
"narHash": "sha256-298x90knBUWX5GHGXh2SKsAKvStjU2ri9UgOGoF79/8=",
"owner": "ghostty-org",
"repo": "ghostty",
"rev": "88b4cd047fa627cdca6781bc7e7dc8b75a2cecb9",
"type": "github"
},
"original": {
"owner": "ghostty-org",
"repo": "ghostty",
"type": "github"
}
},
"git-hooks": {
"inputs": {
"flake-compat": [
"devenv",
"flake-compat"
],
"nixpkgs": [
"devenv",
"nixpkgs"
]
},
"locked": {
"lastModified": 1782908218,
"narHash": "sha256-wLMOrPgVyeF3XmP+qfYcLqnVdTxikdcSvbIY7rA9jTA=",
"owner": "cachix",
"repo": "git-hooks.nix",
"rev": "9f7e99119ece7705299595299f3b031f39356de1",
"type": "github"
},
"original": {
"owner": "cachix",
"repo": "git-hooks.nix",
"type": "github"
}
},
"mk-shell-bin": {
"locked": {
"lastModified": 1677004959,
"narHash": "sha256-/uEkr1UkJrh11vD02aqufCxtbF5YnhRTIKlx5kyvf+I=",
"owner": "rrbutani",
"repo": "nix-mk-shell-bin",
"rev": "ff5d8bd4d68a347be5042e2f16caee391cd75887",
"type": "github"
},
"original": {
"owner": "rrbutani",
"repo": "nix-mk-shell-bin",
"type": "github"
}
},
"nix": {
"inputs": {
"flake-compat": [
"devenv",
"flake-compat"
],
"flake-parts": [
"devenv",
"flake-parts"
],
"git-hooks-nix": [
"devenv",
"git-hooks"
],
"nixpkgs": [
"devenv",
"nixpkgs"
],
"nixpkgs-23-11": [
"devenv"
],
"nixpkgs-regression": [
"devenv"
]
},
"locked": {
"lastModified": 1785349663,
"narHash": "sha256-JSD8lPe5kalvPKx5X+inX8ZZdGLeXkAbd3Jiv7UDf+I=",
"owner": "cachix",
"repo": "nix",
"rev": "f33db89fd6db6edc337d93212f6628ab6d25f407",
"type": "github"
},
"original": {
"owner": "cachix",
"ref": "devenv-2.34",
"repo": "nix",
"type": "github"
}
},
"nix2container": {
"inputs": {
"nixpkgs": [
"nixpkgs"
]
},
"locked": {
"lastModified": 1775487831,
"narHash": "sha256-2lguQpLPQaxpQCJjXhmEEAfabwsAhkP29Z7fgLzHARA=",
"owner": "nlewo",
"repo": "nix2container",
"rev": "76be9608a7f4d6c985d28b0e7be903ae2547df3e",
"type": "github"
},
"original": {
"owner": "nlewo",
"repo": "nix2container",
"type": "github"
}
},
"nixd": {
"inputs": {
"flake-parts": [
"devenv",
"flake-parts"
],
"nixpkgs": [
"devenv",
"nixpkgs"
],
"treefmt-nix": "treefmt-nix"
},
"locked": {
"lastModified": 1783935112,
"narHash": "sha256-IAQ14nteIKXAz4cd75UcZrsHGEVJ7QNrkUwX9rmeZ/Y=",
"owner": "nix-community",
"repo": "nixd",
"rev": "a64cd33e53b316b6b092ea0a966640cd2309bf3d",
"type": "github"
},
"original": {
"owner": "nix-community",
"repo": "nixd",
"type": "github"
}
},
"nixpkgs": {
"locked": {
"lastModified": 1772624091,
"narHash": "sha256-QKyJ0QGWBn6r0invrMAK8dmJoBYWoOWy7lN+UHzW1jc=",
"owner": "NixOS",
"repo": "nixpkgs",
"rev": "80bdc1e5ce51f56b19791b52b2901187931f5353",
"type": "github"
},
"original": {
"owner": "NixOS",
"ref": "nixos-unstable",
"repo": "nixpkgs",
"type": "github"
}
},
"nixpkgs-lib": {
"locked": {
"lastModified": 1782614948,
"narHash": "sha256-ePjCwr1sNm9NYUqywL7QfK3JnlS015msC+eBu2zKlp8=",
"owner": "nix-community",
"repo": "nixpkgs.lib",
"rev": "db3f255737b94216eb71cce308e2912cf6bc2d7c",
"type": "github"
},
"original": {
"owner": "nix-community",
"repo": "nixpkgs.lib",
"type": "github"
}
},
"nixpkgs_2": {
"locked": {
"lastModified": 1785571196,
"narHash": "sha256-KoTsyMQqnXQZq8deCEnu4QkyldkwH/bpMMhUcfMdGIw=",
"owner": "NixOS",
"repo": "nixpkgs",
"rev": "148bab9c1c3c53136ecb44a6ea356a0ed5b39b06",
"type": "github"
},
"original": {
"owner": "NixOS",
"ref": "nixos-unstable",
"repo": "nixpkgs",
"type": "github"
}
},
"root": {
"inputs": {
"devenv": "devenv",
"flake-parts": "flake-parts_2",
"mk-shell-bin": "mk-shell-bin",
"nix2container": "nix2container",
"nixpkgs": "nixpkgs_2",
"rust-overlay": "rust-overlay_2"
}
},
"rust-overlay": {
"inputs": {
"nixpkgs": [
"devenv",
"nixpkgs"
]
},
"locked": {
"lastModified": 1782875958,
"narHash": "sha256-5eqDcnBjb1424HRQdnhuhNOBZguq1Z2tqSa2OMF/m2c=",
"owner": "oxalica",
"repo": "rust-overlay",
"rev": "13139aefa973f3d96c60c0fbab801de058ae25ca",
"type": "github"
},
"original": {
"owner": "oxalica",
"repo": "rust-overlay",
"type": "github"
}
},
"rust-overlay_2": {
"inputs": {
"nixpkgs": [
"nixpkgs"
]
},
"locked": {
"lastModified": 1785562362,
"narHash": "sha256-J15aBa3d6B1SUUAQydQ06wFjPzrrcVceb6jkdfwfGls=",
"owner": "oxalica",
"repo": "rust-overlay",
"rev": "5f29c219a7655519f8a9f8c6968064b82c17cc93",
"type": "github"
},
"original": {
"owner": "oxalica",
"repo": "rust-overlay",
"type": "github"
}
},
"treefmt-nix": {
"inputs": {
"nixpkgs": [
"devenv",
"nixd",
"nixpkgs"
]
},
"locked": {
"lastModified": 1780220602,
"narHash": "sha256-eynAfOmbmxJnkp7YewvCEbShNnnYJ9gLLqkzsYtBPeM=",
"owner": "numtide",
"repo": "treefmt-nix",
"rev": "db947814a175b7ca6ded66e21383d938df01c227",
"type": "github"
},
"original": {
"owner": "numtide",
"repo": "treefmt-nix",
"type": "github"
}
}
},
"root": "root",
"version": 7
}
+33
View File
@@ -0,0 +1,33 @@
{
description = "Furumi TUI development environment";
inputs = {
nixpkgs.url = "github:NixOS/nixpkgs/nixos-unstable";
flake-parts.url = "github:hercules-ci/flake-parts";
devenv.url = "github:cachix/devenv";
devenv.inputs.nixpkgs.follows = "nixpkgs";
rust-overlay.url = "github:oxalica/rust-overlay";
rust-overlay.inputs.nixpkgs.follows = "nixpkgs";
nix2container.url = "github:nlewo/nix2container";
nix2container.inputs.nixpkgs.follows = "nixpkgs";
mk-shell-bin.url = "github:rrbutani/nix-mk-shell-bin";
};
outputs = inputs@{ flake-parts, ... }:
flake-parts.lib.mkFlake { inherit inputs; } {
systems = [
"x86_64-linux"
"aarch64-linux"
"x86_64-darwin"
"aarch64-darwin"
];
imports = [ inputs.devenv.flakeModule ];
perSystem = { ... }: {
devenv.shells.default = {
imports = [ ./devenv.nix ];
};
};
};
}
+4
View File
@@ -0,0 +1,4 @@
[toolchain]
channel = "1.97.0"
profile = "minimal"
components = ["clippy", "rust-src", "rustfmt"]
+6
View File
@@ -38,6 +38,8 @@ pub enum Action {
OpenCurrentTrackInfo,
QueueAddNext,
QueueAddLast,
MoveQueueUp,
MoveQueueDown,
/// Download the selected federated track(s) into the local library.
DownloadSelected,
RemoveFromQueue,
@@ -107,6 +109,8 @@ impl Action {
| Action::OpenListenHistory => Category::Playback,
Action::QueueAddNext
| Action::QueueAddLast
| Action::MoveQueueUp
| Action::MoveQueueDown
| Action::DownloadSelected
| Action::RemoveFromQueue
| Action::ClearQueue
@@ -193,6 +197,8 @@ impl Action {
Action::OpenCurrentTrackInfo => "Current track info".into(),
Action::QueueAddNext => "Queue: add next".into(),
Action::QueueAddLast => "Queue: add to end".into(),
Action::MoveQueueUp => "Queue: move track/selection up".into(),
Action::MoveQueueDown => "Queue: move track/selection down".into(),
Action::DownloadSelected => "Federation: download to library".into(),
Action::RemoveFromQueue => "Queue: remove selected".into(),
Action::ClearQueue => "Queue: clear".into(),
+13
View File
@@ -17,6 +17,16 @@ pub fn handle_key(state: &mut AppState, runtime: &mut Runtime, key: KeyEvent) {
match key.code {
KeyCode::Esc => cancel(state),
KeyCode::Enter => commit(state, runtime),
KeyCode::Up => {
if state.cmdline.history_previous() {
after_change(state, runtime);
}
}
KeyCode::Down => {
if state.cmdline.history_next() {
after_change(state, runtime);
}
}
// Backspace on an empty line closes it, like vim.
KeyCode::Backspace if state.cmdline.input.is_empty() => cancel(state),
_ => {
@@ -162,7 +172,9 @@ pub(super) fn refresh_local_search(state: &mut AppState, runtime: &Runtime) {
/// Enter: close the line. Live commands already took effect (their view
/// stays open); one-shot commands execute here.
fn commit(state: &mut AppState, runtime: &mut Runtime) {
let remembered = state.cmdline.input.as_str().to_string();
let parsed = command::parse(&state.cmdline.input);
state.cmdline.remember(&remembered);
close(state);
match parsed {
Parsed::Empty => {}
@@ -272,6 +284,7 @@ fn close(state: &mut AppState) {
state.cmdline.active = false;
state.cmdline.input.clear();
state.cmdline.live = false;
state.cmdline.begin_history_navigation();
}
/// Pop the live search view if this command-line session opened it.
+2
View File
@@ -56,6 +56,8 @@ pub enum AppEvent {
LocalContentIdsLoaded(Result<Vec<String>, String>),
/// Counts and storage footprint of the local library/database.
LocalLibraryStatsLoaded(Result<crate::library::LocalLibraryStats, String>),
MusicDirectoryValidated(Result<std::path::PathBuf, String>),
MusicDirectoryChanged(Result<crate::library::MusicRelocationStats, String>),
ListenHistoryLoaded(Result<Vec<crate::library::ListenHistoryEntry>, String>),
/// One content id became available locally while the UI is open.
LocalContentAvailable {
+179 -8
View File
@@ -8,6 +8,7 @@ pub mod state;
pub mod update;
use std::io;
use std::io::Write as _;
use std::path::{Path, PathBuf};
use std::process::Command;
use std::sync::Arc;
@@ -61,6 +62,8 @@ pub struct Runtime {
pub art_semaphore: Arc<tokio::sync::Semaphore>,
/// The terminal screen was externally disturbed and needs a full repaint.
pub force_redraw: bool,
/// The alternate screen is temporarily suspended for copy-friendly text.
pub plain_text_mode: bool,
/// Monotonic sequence for live search; stale responses are dropped.
pub search_seq: Arc<std::sync::atomic::AtomicU64>,
pub player: player::Controller,
@@ -106,6 +109,20 @@ fn refresh_local_library_stats(runtime: &Runtime) {
});
}
pub(super) fn validate_music_directory(state: &mut AppState, runtime: &Runtime, path: PathBuf) {
if state.music_dir_changing {
state.status_message = Some("music directory change is already running".into());
return;
}
state.music_dir_changing = true;
state.status_message = Some("checking music directory write access…".into());
let tx = runtime.event_tx.clone();
tokio::task::spawn_blocking(move || {
let result = Library::validate_music_directory(&path).map_err(err_string);
let _ = tx.send(AppEvent::MusicDirectoryValidated(result));
});
}
fn spawn_artist_federation_enrichment(runtime: &Runtime, id: i64, name: String) {
let fed = Arc::clone(&runtime.federation);
let tx = runtime.event_tx.clone();
@@ -244,6 +261,7 @@ pub async fn run(
};
state.player.volume = settings.volume;
state.global.filters = settings.library;
state.music_dir = settings.music_dir.clone();
if let Err(err) = state.visualizer.load_library() {
state.status_message = Some(format!("visualizations disabled: {err:#}"));
}
@@ -255,7 +273,9 @@ pub async fn run(
Arc::clone(&library),
Arc::clone(&devices),
Arc::clone(&jam),
settings.music_dir.clone(),
);
state.music_dir = federation.media_dir();
state.federation.settings = federation.settings();
state.federation.devices = Some(devices.status());
if let Ok((device_id, device_name)) = devices.identity_summary() {
@@ -263,6 +283,7 @@ pub async fn run(
state.device_playback.self_device_name = device_name.clone();
state.device_playback.active_device_id = Some(device_id);
state.device_playback.active_device_name = Some(device_name);
state.device_playback.startup_takeover_pending = true;
}
let player_events = event_tx.clone();
let mut runtime = Runtime {
@@ -276,7 +297,7 @@ pub async fn run(
library_network_refreshing: Arc::new(std::sync::atomic::AtomicBool::new(false)),
library_network_cursors: Arc::new(std::sync::Mutex::new(std::collections::HashMap::new())),
library_network_done: Arc::new(std::sync::Mutex::new(std::collections::HashSet::new())),
library_network_mode: crate::config::settings::LibrarySourceMode::Local,
library_network_mode: state.global.filters.source_mode,
library_network_art_fetching: Arc::new(std::sync::atomic::AtomicBool::new(false)),
library_network_art_attempted: Arc::new(std::sync::Mutex::new(
std::collections::HashSet::new(),
@@ -287,6 +308,7 @@ pub async fn run(
fed_streaming: Arc::new(std::sync::Mutex::new(std::collections::HashSet::new())),
art_semaphore: Arc::new(tokio::sync::Semaphore::new(4)),
force_redraw: false,
plain_text_mode: false,
search_seq: Arc::new(std::sync::atomic::AtomicU64::new(0)),
player: player::spawn(move |event| {
let _ = player_events.send(AppEvent::Player(event));
@@ -317,11 +339,30 @@ pub async fn run(
visual_tick.set_missed_tick_behavior(MissedTickBehavior::Skip);
loop {
if runtime.force_redraw {
terminal.clear()?;
runtime.force_redraw = false;
let plain_text = match state.popup.as_ref() {
Some(state::Popup::PlainText { text }) => Some(text.as_str()),
_ => None,
};
match (runtime.plain_text_mode, plain_text) {
(false, Some(text)) => {
enter_plain_text_mode(text)?;
runtime.plain_text_mode = true;
}
(true, None) => {
leave_plain_text_mode()?;
runtime.plain_text_mode = false;
runtime.force_redraw = true;
}
_ => {}
}
if !runtime.plain_text_mode {
if runtime.force_redraw {
terminal.clear()?;
runtime.force_redraw = false;
}
terminal.draw(|frame| ui::draw(frame, &state, &keymap))?;
}
terminal.draw(|frame| ui::draw(frame, &state, &keymap))?;
tokio::select! {
maybe_event = input.next() => match maybe_event {
@@ -346,6 +387,10 @@ pub async fn run(
}
if state.should_quit {
if runtime.plain_text_mode {
leave_plain_text_mode()?;
runtime.plain_text_mode = false;
}
state.shutting_down = true;
terminal.draw(|frame| ui::draw(frame, &state, &keymap))?;
runtime.federation.shutdown().await;
@@ -355,6 +400,32 @@ pub async fn run(
}
}
fn enter_plain_text_mode(text: &str) -> Result<()> {
let text: String = text
.chars()
.filter(|character| !character.is_control())
.collect();
crossterm::execute!(
io::stdout(),
crossterm::terminal::LeaveAlternateScreen,
crossterm::event::DisableBracketedPaste,
crossterm::terminal::Clear(crossterm::terminal::ClearType::All),
crossterm::cursor::MoveTo(0, 0),
crossterm::style::Print(text)
)?;
io::stdout().flush()?;
Ok(())
}
fn leave_plain_text_mode() -> Result<()> {
crossterm::execute!(
io::stdout(),
crossterm::terminal::EnterAlternateScreen,
crossterm::event::EnableBracketedPaste
)?;
Ok(())
}
fn sync_player_shared(state: &mut AppState, runtime: &Runtime) {
if state.device_playback.is_control() {
extrapolate_control_position(state);
@@ -452,6 +523,7 @@ fn apply_playback_state_to_ui(
.filter(|track| update::track_allowed_by_source_mode(state, track))
.collect();
state.player.queue_pos = queue_pos.min(state.player.queue.len().saturating_sub(1));
state.player.play_next_end = None;
state.player.playing = wire.playing && !state.player.queue.is_empty();
state.player.paused = wire.paused;
state.device_playback.local_idle_since_ms = if state.player.playing && !state.player.paused {
@@ -1330,6 +1402,14 @@ fn perform_effect(state: &mut AppState, runtime: &mut Runtime, effect: Effect) {
}
Effect::SetOptions => {}
Effect::PlaybackQueueChanged => {}
Effect::QueueOrderChanged { restart_current } => {
if restart_current && state.player.playing && state.player.current.is_some() {
let paused = state.player.paused;
start_current_audio(state, runtime, state.player.position_secs, paused);
push_media_metadata(state, runtime);
push_media_update(state, runtime, true);
}
}
Effect::SourceModeChanged => {
runtime.library_network_refresh_at = None;
if let Ok(mut cursors) = runtime.library_network_cursors.lock() {
@@ -1348,6 +1428,30 @@ fn perform_effect(state: &mut AppState, runtime: &mut Runtime, effect: Effect) {
perform_effect(state, runtime, effect);
}
}
Effect::ChangeMusicDirectory {
path,
move_existing,
} => {
if state.music_dir_changing {
state.status_message = Some("music directory change is already running".into());
return;
}
state.music_dir_changing = true;
state.status_message = Some(if move_existing {
"moving saved music to the new directory…".into()
} else {
"changing music save directory…".into()
});
let federation = Arc::clone(&runtime.federation);
let tx = runtime.event_tx.clone();
tokio::spawn(async move {
let result = federation
.change_media_dir(path, move_existing)
.await
.map_err(err_string);
let _ = tx.send(AppEvent::MusicDirectoryChanged(result));
});
}
Effect::EnqueueRelease { id, next } => {
let library = Arc::clone(&runtime.library);
let tx = runtime.event_tx.clone();
@@ -1732,6 +1836,7 @@ fn is_controlled_playback_effect(effect: &Effect) -> bool {
| Effect::SetOptions
| Effect::RemoveQueueIndices { .. }
| Effect::PlaybackQueueChanged
| Effect::QueueOrderChanged { .. }
)
}
@@ -1765,7 +1870,9 @@ fn perform_control_playback_effect(state: &mut AppState, runtime: &mut Runtime,
Effect::SetOptions
| Effect::RemoveQueueIndices { .. }
| Effect::PlaybackQueueChanged
| Effect::LoadListenHistory => {}
| Effect::QueueOrderChanged { .. }
| Effect::LoadListenHistory
| Effect::ChangeMusicDirectory { .. } => {}
_ => {}
}
if local_only_volume {
@@ -2677,6 +2784,7 @@ fn save_app_settings(state: &AppState) {
let settings = crate::config::settings::AppSettings {
volume: state.player.volume,
library: state.global.filters,
music_dir: state.music_dir.clone(),
};
if let Err(err) = crate::config::settings::save(&settings) {
tracing::warn!(%err, "saving app settings failed");
@@ -2811,6 +2919,18 @@ fn handle_device_playback_snapshot(
if !snapshot.active {
return;
}
// Starting a player is an explicit claim of the active role. Import the
// current queue/position from the previously active peer, then announce a
// normal handoff so that the old owner becomes a control device. This is
// intentionally one-shot: subsequent snapshots use the regular lease and
// explicit-transfer rules.
if state.device_playback.startup_takeover_pending {
state.device_playback.startup_takeover_pending = false;
become_control_device(state, runtime, snapshot);
transfer_active_to_this_device(state, runtime);
state.status_message = Some("playback moved to this newly started player".to_string());
return;
}
if local_active_lease_protected(state, now) {
tracing::debug!(
remote = %snapshot.device_id,
@@ -2926,6 +3046,48 @@ fn handle_app_event(state: &mut AppState, runtime: &mut Runtime, event: AppEvent
Err(err) => state::Loadable::Failed(err),
});
}
AppEvent::MusicDirectoryValidated(result) => {
state.music_dir_changing = false;
match result {
Ok(path) => {
let current = std::fs::canonicalize(&state.music_dir)
.unwrap_or_else(|_| state.music_dir.clone());
if path == current {
state.status_message =
Some("this is already the music save directory".into());
} else {
state.popup = Some(state::Popup::ConfirmMusicDirectory { path });
state.status_message = None;
}
}
Err(message) => {
state.status_message = Some(format!(
"music directory is not writable; nothing changed: {message}"
));
}
}
}
AppEvent::MusicDirectoryChanged(result) => {
state.music_dir_changing = false;
match result {
Ok(stats) => {
state.music_dir = runtime.federation.media_dir();
save_app_settings(state);
state.status_message = Some(format!(
"music directory changed · moved {} track(s), {} image(s)",
stats.tracks, stats.images
));
let _ = runtime.event_tx.send(AppEvent::LibraryChanged {
message: state.status_message.clone(),
});
}
Err(message) => {
state.status_message = Some(format!(
"music directory change failed; old library kept: {message}"
));
}
}
}
AppEvent::FederationStatus(status) => {
state.federation.status = Some(status);
}
@@ -3005,6 +3167,7 @@ fn handle_app_event(state: &mut AppState, runtime: &mut Runtime, event: AppEvent
title: "Device invite".to_string(),
text: invite,
help: "Use this invite on another client within 10 minutes to pair it with this device group.".to_string(),
cursor: 0,
});
state.status_message = Some("device invite generated".to_string());
}
@@ -3104,6 +3267,7 @@ fn handle_app_event(state: &mut AppState, runtime: &mut Runtime, event: AppEvent
title: "Jam invite".to_string(),
text: invite,
help: "Copied capability lets federation peers control this host player until restart or regeneration.".to_string(),
cursor: 0,
});
state.status_message = Some("Jam started".into());
}
@@ -3260,6 +3424,7 @@ fn handle_app_event(state: &mut AppState, runtime: &mut Runtime, event: AppEvent
// Skip the failed track instead of stalling the queue.
if state.player.queue_pos + 1 < state.player.queue.len() {
state.player.queue_pos += 1;
update::normalize_play_next_block(&mut state.player);
start_current_audio(state, runtime, 0.0, state.player.paused);
push_media_metadata(state, runtime);
push_media_update(state, runtime, true);
@@ -3361,6 +3526,7 @@ fn handle_app_event(state: &mut AppState, runtime: &mut Runtime, event: AppEvent
text: ticket,
help: "Copy this ticket and paste it into Connect to a peer on another client."
.to_string(),
cursor: 0,
});
}
Err(message) => state.status_message = Some(message),
@@ -3503,6 +3669,7 @@ fn handle_app_event(state: &mut AppState, runtime: &mut Runtime, event: AppEvent
.take()
.unwrap_or(state.player.queue_pos + 1);
state.player.queue_pos = next_pos.min(state.player.queue.len().saturating_sub(1));
update::normalize_play_next_block(&mut state.player);
state.player.current = state.player.queue.get(state.player.queue_pos).cloned();
state.player.position_secs = 0.0;
state.player.track_started_at = Some(now_epoch_seconds());
@@ -3635,9 +3802,13 @@ fn handle_app_event(state: &mut AppState, runtime: &mut Runtime, event: AppEvent
}
AppEvent::EnqueueTracks { tracks, next } => {
let previous_len = state.player.queue.len();
update::enqueue_tracks(state, tracks, next);
let restart_current = update::enqueue_tracks(state, tracks, next);
let count = state.player.queue.len().saturating_sub(previous_len);
record_control_playback_state(state, runtime, false);
perform_effect(
state,
runtime,
Effect::QueueOrderChanged { restart_current },
);
state.status_message = Some(if count == 0 {
"no tracks available in the current source mode".to_string()
} else if next {
+79 -7
View File
@@ -136,17 +136,52 @@ pub fn handle_key(state: &mut AppState, runtime: &mut Runtime, key: KeyEvent) {
KeyCode::Esc | KeyCode::Enter | KeyCode::Char('q') => {}
_ => state.popup = Some(Popup::FedText { title, text }),
},
Popup::FedCopyText { title, text, help } => match key.code {
Popup::FedCopyText {
title,
text,
help,
mut cursor,
} => match key.code {
KeyCode::Esc | KeyCode::Char('q') => {}
KeyCode::Left | KeyCode::Right | KeyCode::Tab | KeyCode::BackTab => {
cursor = usize::from(cursor == 0);
state.popup = Some(Popup::FedCopyText {
title,
text,
help,
cursor,
});
}
KeyCode::Enter if cursor == 1 => state.popup = Some(Popup::PlainText { text }),
KeyCode::Char('p') => state.popup = Some(Popup::PlainText { text }),
KeyCode::Enter | KeyCode::Char('c') => match copy_to_clipboard(&text) {
Ok(()) => state.status_message = Some("copied to clipboard".to_string()),
Err(err) => {
state.status_message = Some(format!("copy failed: {err}"));
state.popup = Some(Popup::FedCopyText { title, text, help });
state.popup = Some(Popup::FedCopyText {
title,
text,
help,
cursor,
});
}
},
_ => state.popup = Some(Popup::FedCopyText { title, text, help }),
_ => {
state.popup = Some(Popup::FedCopyText {
title,
text,
help,
cursor,
})
}
},
Popup::PlainText { text } => match key.code {
KeyCode::Esc => {}
_ => state.popup = Some(Popup::PlainText { text }),
},
Popup::ConfirmMusicDirectory { path } => {
handle_music_directory_confirmation(state, runtime, path, key)
}
Popup::FederationStatusDetails {
focus,
status_cursor,
@@ -502,6 +537,13 @@ fn handle_fed_input(
!state.federation.settings.network_id.is_empty();
super::fed_apply_settings(state, runtime);
}
FedInputField::MusicDirectory => {
if value.is_empty() {
state.status_message = Some("music directory is empty".into());
} else {
super::validate_music_directory(state, runtime, value.into());
}
}
FedInputField::ConnectTicket => {
if value.is_empty() {
state.status_message = Some("ticket is empty".into());
@@ -551,6 +593,31 @@ fn handle_fed_input(
}
}
fn handle_music_directory_confirmation(
state: &mut AppState,
runtime: &mut Runtime,
path: std::path::PathBuf,
key: KeyEvent,
) {
let move_existing = match key.code {
KeyCode::Esc | KeyCode::Char('q') => return,
KeyCode::Enter | KeyCode::Char('y') | KeyCode::Char('Y') => true,
KeyCode::Char('n') | KeyCode::Char('N') => false,
_ => {
state.popup = Some(Popup::ConfirmMusicDirectory { path });
return;
}
};
super::perform_effect(
state,
runtime,
crate::app::update::Effect::ChangeMusicDirectory {
path,
move_existing,
},
);
}
fn handle_device_pairing(
state: &mut AppState,
runtime: &Runtime,
@@ -934,10 +1001,15 @@ fn handle_track_info(
KeyCode::Char('c') => {
if let Some(track) = tracks.get(cursor.min(len.saturating_sub(1))) {
match crate::share::track_share_link(track) {
Some(link) => match copy_to_clipboard(&link) {
Ok(()) => state.status_message = Some("frid link copied".into()),
Err(err) => state.status_message = Some(format!("copy failed: {err}")),
},
Some(link) => {
state.popup = Some(Popup::FedCopyText {
title: "Track share link".into(),
text: link,
help: "Copy the link with a clipboard helper, or show it as one clean terminal line for mouse selection.".into(),
cursor: 0,
});
return;
}
None => state.status_message = Some("no content id for this track yet".into()),
}
}
+119 -1
View File
@@ -389,6 +389,7 @@ mod tests {
featured_tracks: Vec::new(),
};
let mut state = AppState::default();
state.global.filters.source_mode = crate::config::settings::LibrarySourceMode::Local;
state.artist_fed_views.insert(
detail.id,
Loadable::Ready(crate::federation::FedArtistCard {
@@ -716,7 +717,16 @@ pub enum Popup {
title: String,
text: String,
help: String,
cursor: usize,
},
/// A copy-friendly terminal screen containing exactly one logical line.
/// The app loop temporarily leaves the alternate screen while this is
/// active so terminal selection does not acquire TUI borders or hard
/// line breaks.
PlainText { text: String },
/// The destination is already write-tested; Enter/y migrates managed
/// files, while n changes only the destination for future downloads.
ConfirmMusicDirectory { path: std::path::PathBuf },
/// Incoming trusted-device pairing request.
DevicePairing {
request_id: String,
@@ -801,6 +811,7 @@ impl StatusDetailFocus {
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum FedInputField {
MusicDirectory,
NetworkId,
ConnectTicket,
DeviceName,
@@ -811,6 +822,7 @@ pub enum FedInputField {
impl FedInputField {
pub fn title(self) -> &'static str {
match self {
FedInputField::MusicDirectory => "Music save directory",
FedInputField::NetworkId => "Network ID",
FedInputField::ConnectTicket => "Connect to peer (paste ticket)",
FedInputField::DeviceName => "Device name",
@@ -821,6 +833,9 @@ impl FedInputField {
pub fn help(self) -> &'static str {
match self {
FedInputField::MusicDirectory => {
"Federated tracks saved to your library use this directory. The directory is checked for write access before anything changes."
}
FedInputField::NetworkId => {
"A unique network id. It must match exactly on every client that should see and connect to the same peers."
}
@@ -867,6 +882,7 @@ impl FedRow {
/// the config directory.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum SettingsRow {
MusicDirectory,
Federation(FedRow),
StatusDetails,
DeviceName,
@@ -1004,7 +1020,7 @@ pub fn device_status_order(state: &AppState) -> Vec<usize> {
}
pub fn settings_rows(state: &AppState) -> Vec<SettingsRow> {
let mut rows = Vec::new();
let mut rows = vec![SettingsRow::MusicDirectory];
rows.extend(FedRow::ALL.into_iter().map(SettingsRow::Federation));
rows.push(SettingsRow::DeviceName);
rows.push(SettingsRow::DeviceInvite);
@@ -1064,6 +1080,98 @@ pub struct Cmdline {
/// A live command (search) applied effects during this session; Esc
/// undoes them, Enter keeps them.
pub live: bool,
/// Commands committed during this process, oldest first.
pub history: Vec<String>,
/// Entry currently recalled with Up/Down. `None` is the editable draft
/// after the newest history entry.
history_index: Option<usize>,
history_draft: String,
}
impl Cmdline {
const HISTORY_LIMIT: usize = 100;
pub fn begin_history_navigation(&mut self) {
self.history_index = None;
self.history_draft.clear();
}
pub fn remember(&mut self, value: &str) {
let value = value.trim();
if value.is_empty() {
return;
}
if self.history.last().is_none_or(|last| last != value) {
self.history.push(value.to_string());
if self.history.len() > Self::HISTORY_LIMIT {
self.history.remove(0);
}
}
self.begin_history_navigation();
}
pub fn history_previous(&mut self) -> bool {
if self.history.is_empty() {
return false;
}
let index = match self.history_index {
Some(index) => index.saturating_sub(1),
None => {
self.history_draft = self.input.as_str().to_string();
self.history.len() - 1
}
};
self.history_index = Some(index);
self.input = LineEdit::new(self.history[index].clone());
true
}
pub fn history_next(&mut self) -> bool {
let Some(index) = self.history_index else {
return false;
};
if index + 1 < self.history.len() {
self.history_index = Some(index + 1);
self.input = LineEdit::new(self.history[index + 1].clone());
} else {
self.history_index = None;
self.input = LineEdit::new(std::mem::take(&mut self.history_draft));
}
true
}
}
#[cfg(test)]
mod cmdline_history_tests {
use super::Cmdline;
use crate::app::input::LineEdit;
#[test]
fn history_walks_oldest_and_restores_the_draft() {
let mut cmdline = Cmdline::default();
cmdline.remember("volume 20");
cmdline.remember("/ambient");
cmdline.input = LineEdit::new("unfinished");
assert!(cmdline.history_previous());
assert_eq!(cmdline.input.as_str(), "/ambient");
assert!(cmdline.history_previous());
assert_eq!(cmdline.input.as_str(), "volume 20");
assert!(cmdline.history_previous());
assert_eq!(cmdline.input.as_str(), "volume 20");
assert!(cmdline.history_next());
assert_eq!(cmdline.input.as_str(), "/ambient");
assert!(cmdline.history_next());
assert_eq!(cmdline.input.as_str(), "unfinished");
}
#[test]
fn history_deduplicates_consecutive_commands() {
let mut cmdline = Cmdline::default();
cmdline.remember("q");
cmdline.remember(" q ");
assert_eq!(cmdline.history, vec!["q"]);
}
}
/// Live search state driven by the `:/query` command.
@@ -1169,6 +1277,9 @@ impl RepeatMode {
pub struct PlayerBar {
pub queue: Vec<TrackItem>,
pub queue_pos: usize,
/// Exclusive end of the contiguous "play next" block built by
/// sequential `a` actions. New `a` additions are inserted here.
pub play_next_end: Option<usize>,
pub current: Option<TrackItem>,
/// A track is loaded (playing or paused); false = stopped.
pub playing: bool,
@@ -1194,6 +1305,7 @@ impl Default for PlayerBar {
Self {
queue: Vec::new(),
queue_pos: 0,
play_next_end: None,
current: None,
playing: false,
paused: false,
@@ -1240,6 +1352,9 @@ pub struct DevicePlaybackState {
pub remote: BTreeMap<String, crate::devices::PlaybackSnapshot>,
pub last_remote_snapshot: Option<crate::devices::PlaybackSnapshot>,
pub jam_host: bool,
/// A freshly started TUI owns playback by protocol. The first active
/// snapshot discovered during startup is imported and handed off here.
pub startup_takeover_pending: bool,
}
impl DevicePlaybackState {
@@ -1280,6 +1395,9 @@ pub struct AppState {
pub status_message: Option<String>,
pub spinner_frame: usize,
pub settings_cursor: usize,
/// Root for music permanently downloaded from federation peers.
pub music_dir: std::path::PathBuf,
pub music_dir_changing: bool,
pub player: PlayerBar,
pub device_playback: DevicePlaybackState,
pub jam: crate::jam::JamStatus,
+131 -11
View File
@@ -48,8 +48,19 @@ pub enum Effect {
},
/// Queue/options changed without a direct audio engine action.
PlaybackQueueChanged,
/// Queue order changed. A decoded gapless-next source must be discarded
/// by restarting the current source at its current position.
QueueOrderChanged {
restart_current: bool,
},
/// Persist and apply a Local / My / Global source-mode change.
SourceModeChanged,
/// Switch the permanent federation-download root, optionally relocating
/// files managed under the previous root first.
ChangeMusicDirectory {
path: std::path::PathBuf,
move_existing: bool,
},
/// Persist the federation settings and start/stop the node.
FedApplySettings,
/// Force an immediate library publish into the DHT.
@@ -180,6 +191,7 @@ pub fn update(state: &mut AppState, action: Action) -> Option<Effect> {
}
Action::ToggleShuffle => {
state.player.shuffle = !state.player.shuffle;
state.player.play_next_end = None;
// Shuffle physically reorders the unplayed tail, so the Queue
// tab always shows the real upcoming order; turning it off
// restores the original ordering.
@@ -275,6 +287,7 @@ pub fn update(state: &mut AppState, action: Action) -> Option<Effect> {
Action::OpenCommandLine => {
state.cmdline.active = true;
state.cmdline.input.clear();
state.cmdline.begin_history_navigation();
None
}
Action::OpenSearch => {
@@ -283,6 +296,7 @@ pub fn update(state: &mut AppState, action: Action) -> Option<Effect> {
state.cmdline.active = true;
state.cmdline.input = crate::app::input::LineEdit::new("/");
state.cmdline.live = true;
state.cmdline.begin_history_navigation();
state.search = SearchState::default();
state.active_tab = Tab::Global;
if !matches!(state.global.stack.last(), Some(GlobalView::Search { .. })) {
@@ -371,6 +385,8 @@ pub fn update(state: &mut AppState, action: Action) -> Option<Effect> {
Action::RemoveFromQueue => remove_selected_from_queue(state),
Action::QueueAddNext => queue_add(state, true),
Action::QueueAddLast => queue_add(state, false),
Action::MoveQueueUp => move_selected_queue(state, true),
Action::MoveQueueDown => move_selected_queue(state, false),
Action::GoToRelease => {
let track = selected_track(state).or_else(|| state.player.current.clone());
match track {
@@ -418,6 +434,7 @@ pub fn update(state: &mut AppState, action: Action) -> Option<Effect> {
let had_tracks = !state.player.queue.is_empty();
state.player.queue.clear();
state.player.queue_pos = 0;
state.player.play_next_end = None;
state.player.current = None;
state.player.playing = false;
state.player.paused = false;
@@ -1030,6 +1047,68 @@ fn selected_queue_indices(state: &AppState) -> Vec<usize> {
.unwrap_or_else(|| vec![state.queue_tab.cursor.min(state.player.queue.len() - 1)])
}
fn move_selected_queue(state: &mut AppState, up: bool) -> Option<Effect> {
let indices = selected_queue_indices(state);
let Some(&start) = indices.first() else {
state.status_message = Some("queue is empty".into());
return None;
};
let end = *indices.last().unwrap_or(&start);
let len = state.player.queue.len();
if (up && start == 0) || (!up && end + 1 >= len) {
state.status_message = Some(if up {
"selection is already at the top of the queue".into()
} else {
"selection is already at the bottom of the queue".into()
});
return None;
}
let old_current = state.player.queue_pos;
let new_current = if up {
let displaced = state.player.queue.remove(start - 1);
state.player.queue.insert(end, displaced);
if old_current == start - 1 {
end
} else if (start..=end).contains(&old_current) {
old_current - 1
} else {
old_current
}
} else {
let displaced = state.player.queue.remove(end + 1);
state.player.queue.insert(start, displaced);
if old_current == end + 1 {
start
} else if (start..=end).contains(&old_current) {
old_current + 1
} else {
old_current
}
};
let delta: isize = if up { -1 } else { 1 };
state.player.queue_pos = new_current;
state.queue_tab.cursor = (state.queue_tab.cursor as isize + delta) as usize;
if state
.track_selection
.is_active_for(&TrackSelectionScope::Queue)
{
state.track_selection.anchor = (state.track_selection.anchor as isize + delta) as usize;
state.track_selection.cursor = (state.track_selection.cursor as isize + delta) as usize;
}
// An explicit manual order supersedes both the temporary play-next block
// and a saved pre-shuffle order.
state.player.play_next_end = None;
state.player.original_order = None;
let restart_current = state.player.prefetched_pos.take().is_some();
state.status_message = Some(format!(
"moved {} track(s) {}",
indices.len(),
if up { "up" } else { "down" }
));
Some(Effect::QueueOrderChanged { restart_current })
}
fn remove_selected_from_queue(state: &mut AppState) -> Option<Effect> {
let indices = selected_queue_indices(state);
if indices.is_empty() {
@@ -1081,6 +1160,10 @@ fn remove_queue_indices(state: &mut AppState, indices: &[usize]) -> QueueRemoval
.iter()
.filter(|index| **index < old_queue_pos)
.count();
if let Some(end) = state.player.play_next_end {
let removed_before_end = unique.iter().filter(|index| **index < end).count();
state.player.play_next_end = Some(end.saturating_sub(removed_before_end));
}
let was_loaded = state.player.playing;
let was_paused = state.player.paused;
@@ -1107,6 +1190,7 @@ fn remove_queue_indices(state: &mut AppState, indices: &[usize]) -> QueueRemoval
state.player.track_started_at = None;
state.player.listen_id = None;
state.queue_tab.cursor = state.queue_tab.cursor.min(state.player.queue.len() - 1);
normalize_play_next_block(&mut state.player);
return QueueRemovalOutcome {
restart_paused: was_loaded.then_some(was_paused),
stop: false,
@@ -1134,6 +1218,7 @@ fn remove_queue_indices(state: &mut AppState, indices: &[usize]) -> QueueRemoval
.cloned()
});
state.queue_tab.cursor = state.queue_tab.cursor.min(state.player.queue.len() - 1);
normalize_play_next_block(&mut state.player);
QueueRemovalOutcome {
restart_paused: None,
stop: false,
@@ -1267,7 +1352,7 @@ fn queue_add(state: &mut AppState, next: bool) -> Option<Effect> {
if !tracks.is_empty() {
let count = tracks.len();
let title = tracks[0].title.clone();
enqueue_tracks(state, tracks, next);
let restart_current = enqueue_tracks(state, tracks, next);
state.track_selection.clear();
state.status_message = Some(if count == 1 && next {
format!("queued next: {title}")
@@ -1278,7 +1363,7 @@ fn queue_add(state: &mut AppState, next: bool) -> Option<Effect> {
} else {
format!("queued: {count} tracks")
});
return Some(Effect::PlaybackQueueChanged);
return Some(Effect::QueueOrderChanged { restart_current });
}
if let Some(id) = selected_release_id(state) {
return Some(Effect::EnqueueRelease { id, next });
@@ -1372,19 +1457,22 @@ pub(crate) fn track_artist_refs(track: &TrackItem) -> Vec<crate::library::models
refs
}
/// Insert tracks after the playing one (`next`) or at the end. Keeps the
/// gapless prefetch index pointing at the same track if items shift.
pub fn enqueue_tracks(state: &mut AppState, tracks: Vec<TrackItem>, next: bool) {
/// Insert tracks into the stable "play next" block (`next`) or at the end.
/// Returns whether a decoded gapless-next source became stale.
pub fn enqueue_tracks(state: &mut AppState, tracks: Vec<TrackItem>, next: bool) -> bool {
let tracks: Vec<_> = tracks
.into_iter()
.filter(|track| track_allowed_by_source_mode(state, track))
.collect();
let player = &mut state.player;
if tracks.is_empty() {
return;
return false;
}
let insert_at = if next && !player.queue.is_empty() {
(player.queue_pos + 1).min(player.queue.len())
player
.play_next_end
.filter(|end| *end >= player.queue_pos.saturating_add(1) && *end <= player.queue.len())
.unwrap_or_else(|| (player.queue_pos + 1).min(player.queue.len()))
} else if next {
0
} else {
@@ -1394,14 +1482,28 @@ pub fn enqueue_tracks(state: &mut AppState, tracks: Vec<TrackItem>, next: bool)
for (offset, track) in tracks.into_iter().enumerate() {
player.queue.insert(insert_at + offset, track);
}
if let Some(prefetched) = &mut player.prefetched_pos
&& insert_at <= *prefetched
{
*prefetched += count;
let restart_current = player
.prefetched_pos
.is_some_and(|prefetched| insert_at <= prefetched);
if restart_current {
player.prefetched_pos = None;
}
if insert_at <= player.queue_pos && player.current.is_some() {
player.queue_pos += count;
}
if next {
player.play_next_end = Some(insert_at + count);
}
restart_current
}
pub(super) fn normalize_play_next_block(player: &mut super::state::PlayerBar) {
if player
.play_next_end
.is_some_and(|end| end <= player.queue_pos || end > player.queue.len())
{
player.play_next_end = None;
}
}
/// Manual queue navigation (n / p); the tail is pre-shuffled when shuffle
@@ -1426,6 +1528,7 @@ fn queue_step(state: &mut AppState, direction: isize) -> Option<Effect> {
} else {
player.queue_pos = next as usize;
}
normalize_play_next_block(player);
Some(Effect::PlayCurrent)
}
@@ -1464,9 +1567,11 @@ pub fn advance_after_finish(state: &mut AppState) -> Option<Effect> {
repeat => {
if player.queue_pos + 1 < player.queue.len() {
player.queue_pos += 1;
normalize_play_next_block(player);
Some(Effect::PlayCurrent)
} else if repeat == super::state::RepeatMode::All {
player.queue_pos = 0;
player.play_next_end = None;
Some(Effect::PlayCurrent)
} else {
player.playing = false;
@@ -2076,6 +2181,7 @@ fn select_current(state: &mut AppState) -> Option<Effect> {
return None;
}
state.player.queue_pos = state.queue_tab.cursor.min(state.player.queue.len() - 1);
normalize_play_next_block(&mut state.player);
return Some(Effect::PlayCurrent);
}
if state.active_tab != Tab::Global {
@@ -2616,6 +2722,19 @@ fn fed_card_featured_artist_names(track: &crate::federation::FedCardTrack) -> Ve
fn federation_select(state: &mut AppState) -> Option<Effect> {
use super::state::{FedInputField, FedRow, Popup, SettingsRow};
match settings_rows(state).get(state.settings_cursor).copied()? {
SettingsRow::MusicDirectory => {
if state.music_dir_changing {
state.status_message = Some("music directory change is already running".into());
return None;
}
state.popup = Some(Popup::FedInput {
field: FedInputField::MusicDirectory,
input: crate::app::input::LineEdit::new(
state.music_dir.to_string_lossy().into_owned(),
),
});
None
}
SettingsRow::Federation(FedRow::Toggle) => {
let settings = &mut state.federation.settings;
if !settings.enabled && settings.network_id.trim().is_empty() {
@@ -2789,6 +2908,7 @@ fn require_connected_devices_enabled(state: &mut AppState) -> bool {
pub(super) fn on_new_queue(state: &mut AppState) {
let player = &mut state.player;
player.original_order = None;
player.play_next_end = None;
if player.shuffle && !player.queue.is_empty() {
player.original_order = Some(player.queue.iter().map(track_key).collect());
shuffle_range(player, (player.queue_pos + 1).min(player.queue.len()));
+72 -5
View File
@@ -158,25 +158,25 @@ fn source_mode_cycles_on_library_playlists_and_queue_tabs() {
update(&mut state, Action::CycleSourceMode),
Some(Effect::SourceModeChanged)
);
assert_eq!(state.global.filters.source_mode, LibrarySourceMode::My);
assert_eq!(state.global.filters.source_mode, LibrarySourceMode::Local);
state.active_tab = Tab::Playlists;
assert_eq!(
update(&mut state, Action::CycleSourceMode),
Some(Effect::SourceModeChanged)
);
assert_eq!(state.global.filters.source_mode, LibrarySourceMode::Global);
assert_eq!(state.global.filters.source_mode, LibrarySourceMode::My);
state.active_tab = Tab::Queue;
assert_eq!(
update(&mut state, Action::CycleSourceMode),
Some(Effect::SourceModeChanged)
);
assert_eq!(state.global.filters.source_mode, LibrarySourceMode::Local);
assert_eq!(state.global.filters.source_mode, LibrarySourceMode::Global);
state.active_tab = Tab::Federation;
assert_eq!(update(&mut state, Action::CycleSourceMode), None);
assert_eq!(state.global.filters.source_mode, LibrarySourceMode::Local);
assert_eq!(state.global.filters.source_mode, LibrarySourceMode::Global);
}
#[test]
@@ -457,6 +457,7 @@ fn local_mode_hides_pending_federation_tracks_from_playlists_and_playback() {
active_tab: Tab::Playlists,
..AppState::default()
};
state.global.filters.source_mode = crate::config::settings::LibrarySourceMode::Local;
state.playlists.opened = Some(OpenedPlaylist { id: 7, cursor: 1 });
state.playlist_views.insert(
7,
@@ -512,6 +513,7 @@ fn network_modes_show_pending_federation_playlist_tracks() {
#[test]
fn local_mode_rejects_async_federation_queue_additions() {
let mut state = AppState::default();
state.global.filters.source_mode = crate::config::settings::LibrarySourceMode::Local;
enqueue_tracks(
&mut state,
@@ -732,13 +734,78 @@ fn artist_top_track_selection_queues_all_selected_tracks() {
update(&mut state, Action::MoveDown);
assert_eq!(
update(&mut state, Action::QueueAddLast),
Some(Effect::PlaybackQueueChanged),
Some(Effect::QueueOrderChanged {
restart_current: false,
}),
);
let queued: Vec<i64> = state.player.queue.iter().map(|track| track.id).collect();
assert_eq!(queued, vec![1, 2]);
assert!(!state.track_selection.is_active());
}
#[test]
fn sequential_queue_next_additions_keep_their_order_as_one_block() {
let mut state = AppState::default();
state.player.queue = (10..=13).map(test_track).collect();
state.player.queue_pos = 0;
state.player.current = Some(test_track(10));
assert!(!enqueue_tracks(&mut state, vec![test_track(1)], true));
assert!(!enqueue_tracks(&mut state, vec![test_track(2)], true));
assert!(!enqueue_tracks(&mut state, vec![test_track(3)], true));
assert_eq!(
state
.player
.queue
.iter()
.map(|track| track.id)
.collect::<Vec<_>>(),
vec![10, 1, 2, 3, 11, 12, 13]
);
assert_eq!(state.player.play_next_end, Some(4));
}
#[test]
fn queue_selection_moves_as_a_group_and_preserves_current_track() {
let mut state = AppState {
active_tab: Tab::Queue,
..AppState::default()
};
state.player.queue = (1..=5).map(test_track).collect();
state.player.queue_pos = 0;
state.player.current = Some(test_track(1));
state.queue_tab.cursor = 2;
state.track_selection.start(TrackSelectionScope::Queue, 2);
state
.track_selection
.set_cursor(TrackSelectionScope::Queue, 3);
assert_eq!(
update(&mut state, Action::MoveQueueUp),
Some(Effect::QueueOrderChanged {
restart_current: false,
})
);
assert_eq!(
state
.player
.queue
.iter()
.map(|track| track.id)
.collect::<Vec<_>>(),
vec![1, 3, 4, 2, 5]
);
assert_eq!(state.player.queue_pos, 0);
assert_eq!(state.queue_tab.cursor, 1);
assert_eq!(
state
.track_selection
.indices(&TrackSelectionScope::Queue, 5),
Some(vec![1, 2])
);
}
#[test]
fn removing_current_queue_track_requests_paused_restart() {
let mut state = AppState {
+10
View File
@@ -63,6 +63,16 @@ command = "QueueAddNext"
key_sequence = "shift-a"
command = "QueueAddLast"
[[keymaps]]
key_sequence = "alt-k"
command = "MoveQueueUp"
context = "queue"
[[keymaps]]
key_sequence = "alt-j"
command = "MoveQueueDown"
context = "queue"
[[keymaps]]
key_sequence = "shift-c"
command = "OpenConnectedDevices"
+19 -2
View File
@@ -1,12 +1,13 @@
use anyhow::{Context as _, Result};
use serde::{Deserialize, Serialize};
use std::path::PathBuf;
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum LibrarySourceMode {
#[default]
Local,
My,
#[default]
Global,
}
@@ -50,6 +51,9 @@ pub struct AppSettings {
pub volume: u8,
#[serde(default)]
pub library: LibraryFilters,
/// Root used for music materialized from federation peers.
#[serde(default = "default_music_dir")]
pub music_dir: PathBuf,
}
impl Default for AppSettings {
@@ -57,6 +61,7 @@ impl Default for AppSettings {
Self {
volume: default_volume(),
library: LibraryFilters::default(),
music_dir: default_music_dir(),
}
}
}
@@ -64,6 +69,9 @@ impl Default for AppSettings {
impl AppSettings {
pub fn normalized(mut self) -> Self {
self.volume = self.volume.min(100);
if self.music_dir.as_os_str().is_empty() {
self.music_dir = default_music_dir();
}
self
}
}
@@ -72,6 +80,14 @@ fn default_volume() -> u8 {
80
}
/// The historical permanent-download location, kept as the default for
/// backward compatibility with existing installations.
pub fn default_music_dir() -> PathBuf {
crate::config::project_dirs()
.map(|dirs| dirs.data_dir().join("federation-media"))
.unwrap_or_else(|| PathBuf::from("federation-media"))
}
pub fn load() -> (AppSettings, Option<String>) {
let Some(path) = settings_path() else {
return (AppSettings::default(), None);
@@ -128,6 +144,7 @@ hide_featured_only = true
assert_eq!(settings.volume, 100);
assert!(settings.library.hide_featured_only);
assert_eq!(settings.library.source_mode, LibrarySourceMode::Local);
assert_eq!(settings.library.source_mode, LibrarySourceMode::Global);
assert_eq!(settings.music_dir, default_music_dir());
}
}
+60 -23
View File
@@ -24,13 +24,16 @@ use crate::library::models::{ArtistRef, TrackItem};
pub const SYNC_ALPN: &[u8] = b"furumi/sync/2";
const CLIENT_VERSION: &str = env!("CARGO_PKG_VERSION");
const PROTOCOL_VERSION: u16 = 2;
pub const PROTOCOL_VERSION: u16 = 2;
const INVITE_TTL_MS: i64 = 10 * 60 * 1000;
const PAIRING_WAIT_MS: i64 = 5 * 60 * 1000;
const PAIRING_RETRY_DELAY: Duration = Duration::from_secs(1);
const RESPONSE_DRAIN_TIMEOUT: Duration = Duration::from_secs(2);
const SYNC_INTERVAL: Duration = Duration::from_secs(2);
const MAX_LINE: usize = 8 * 1024 * 1024;
/// Playback control is ephemeral. Keeping old full-queue commands in a new
/// peer's catch-up batch can make the initial pairing frame arbitrarily large.
const PLAYBACK_COMMAND_TTL_MS: i64 = 5 * 60 * 1_000;
const MAX_OPS_PER_BATCH: usize = 1000;
#[derive(Debug, Clone, PartialEq, Eq)]
@@ -1249,8 +1252,7 @@ impl DeviceSync {
device: &StoredDevice,
transport_stats: Arc<crate::federation::TransportStats>,
) -> Result<()> {
let ticket: PeerTicket = device.endpoint_ticket.parse()?;
let peer = service.connect(ticket).await?;
let peer = resolve_device_peer(&service, device).await?;
let own_ticket = service.ticket().await?.to_string();
let identity = self.ensure_identity()?;
let profile = self.own_profile(&own_ticket)?;
@@ -2174,25 +2176,39 @@ impl DeviceSync {
LEFT JOIN sync_peer_acks a
ON a.peer_device_id = ?1 AND a.origin_device_id = o.origin_device_id
WHERE o.seq > COALESCE(a.max_seq, 0)
ORDER BY o.hlc_ms, o.op_id
LIMIT ?2",
)?;
let rows = stmt.query_map(params![peer_device_id, MAX_OPS_PER_BATCH as i64], |row| {
let payload_json: String = row.get(4)?;
Ok(SyncOpWire {
op_id: row.get(0)?,
origin_device_id: row.get(1)?,
seq: row.get(2)?,
hlc_ms: row.get(3)?,
payload: serde_json::from_str(&payload_json).map_err(|err| {
rusqlite::Error::FromSqlConversionFailure(
4,
rusqlite::types::Type::Text,
Box::new(err),
AND (
o.kind != 'playback_command'
OR (
o.hlc_ms >= ?2
AND json_extract(o.payload_json, '$.target_device_id') = ?1
)
})?,
})
})?;
)
ORDER BY o.hlc_ms, o.op_id
LIMIT ?3",
)?;
let rows = stmt.query_map(
params![
peer_device_id,
now_ms().saturating_sub(PLAYBACK_COMMAND_TTL_MS),
MAX_OPS_PER_BATCH as i64
],
|row| {
let payload_json: String = row.get(4)?;
Ok(SyncOpWire {
op_id: row.get(0)?,
origin_device_id: row.get(1)?,
seq: row.get(2)?,
hlc_ms: row.get(3)?,
payload: serde_json::from_str(&payload_json).map_err(|err| {
rusqlite::Error::FromSqlConversionFailure(
4,
rusqlite::types::Type::Text,
Box::new(err),
)
})?,
})
},
)?;
Ok(rows.collect::<rusqlite::Result<Vec<_>>>()?)
}
@@ -2380,7 +2396,7 @@ impl DeviceSync {
let own = self.ensure_identity()?.device_id;
let conn = lock(&self.conn);
let mut stmt = conn.prepare(
"SELECT device_id, endpoint_ticket
"SELECT device_id, endpoint_id, endpoint_ticket
FROM sync_devices
WHERE trusted_at_ms IS NOT NULL
AND revoked_at_ms IS NULL
@@ -2390,7 +2406,8 @@ impl DeviceSync {
let rows = stmt.query_map([own], |row| {
Ok(StoredDevice {
device_id: row.get(0)?,
endpoint_ticket: row.get(1)?,
endpoint_id: row.get(1)?,
endpoint_ticket: row.get(2)?,
})
})?;
Ok(rows.collect::<rusqlite::Result<Vec<_>>>()?)
@@ -2642,9 +2659,29 @@ impl DeviceSync {
#[derive(Debug)]
struct StoredDevice {
device_id: String,
endpoint_id: String,
endpoint_ticket: String,
}
async fn resolve_device_peer(
service: &MusicDhtService,
device: &StoredDevice,
) -> Result<music_dht::EndpointId> {
if let Ok(peer) = device.endpoint_id.parse::<music_dht::EndpointId>()
&& (service.connected_peers().contains(&peer)
|| service
.known_peers()
.iter()
.any(|contact| contact.peer_id == peer))
{
// The live DHT contact carries a ticket for the current schema. The
// persisted device ticket may predate a schema upgrade.
return Ok(peer);
}
let ticket: PeerTicket = device.endpoint_ticket.parse()?;
service.connect(ticket).await.map_err(Into::into)
}
pub async fn serve_peers(
mut acceptor: StreamAcceptor,
sync: Arc<DeviceSync>,
+52
View File
@@ -228,6 +228,58 @@ fn playback_command_is_targeted_and_deduplicated() {
assert!(rx.try_recv().is_err());
}
#[test]
fn playback_commands_are_caught_up_only_by_their_target_while_fresh() {
let sync = test_sync();
let command = PlaybackCommand::SetState {
state: PlaybackStateWire {
queue: Vec::new(),
queue_pos: 0,
playing: false,
paused: false,
idle_since_ms: None,
position_secs: 0.0,
volume: 42,
shuffle: false,
repeat: PlaybackRepeat::Off,
},
seek: false,
};
sync.record_playback_command("dev_target", command.clone())
.unwrap();
sync.record_playback_command("dev_other", command).unwrap();
let target_ops = sync.ops_for_peer("dev_target").unwrap();
assert_eq!(
target_ops
.iter()
.filter(|op| matches!(op.payload, SyncOpPayload::PlaybackCommand { .. }))
.count(),
1
);
assert!(
sync.ops_for_peer("dev_unknown")
.unwrap()
.iter()
.all(|op| !matches!(op.payload, SyncOpPayload::PlaybackCommand { .. }))
);
lock(&sync.conn)
.execute(
"UPDATE sync_ops
SET hlc_ms = ?1
WHERE kind = 'playback_command'",
[now_ms().saturating_sub(PLAYBACK_COMMAND_TTL_MS + 1)],
)
.unwrap();
assert!(
sync.ops_for_peer("dev_target")
.unwrap()
.iter()
.all(|op| !matches!(op.payload, SyncOpPayload::PlaybackCommand { .. }))
);
}
#[test]
fn newer_device_trust_reactivates_revoked_device() {
let sync = test_sync();
+2
View File
@@ -18,6 +18,8 @@ use crate::library::Library;
/// ALPN of the audio streaming protocol (shared with furumi-fd).
pub const AUDIO_ALPN: &[u8] = b"furumi-fd/audio/1";
/// Version of the audio transfer stream protocol.
pub const AUDIO_PROTOCOL_VERSION: u16 = 1;
/// Maximum size of a JSON protocol line (request or response header).
const MAX_PROTOCOL_LINE: usize = 4096;
+190
View File
@@ -0,0 +1,190 @@
//! Informational publication and observation of protocol versions.
use std::collections::BTreeMap;
use std::sync::{Arc, Mutex, MutexGuard};
use std::time::Duration;
use anyhow::{Context, Result};
pub use music_dht::capabilities::CAPABILITIES_ALPN;
use music_dht::capabilities::{
CAPABILITIES_PROTOCOL_VERSION, CapabilityManifest, CapabilityMessage, read_message,
write_message,
};
use music_dht::{ByteStream, EndpointId, MusicDhtService, StreamAcceptor};
const PROBE_INTERVAL: Duration = Duration::from_secs(30);
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct ProtocolVersions {
pub local: BTreeMap<String, u16>,
pub observed: BTreeMap<String, u16>,
pub observed_peers: usize,
pub newer: Vec<NewerProtocol>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct NewerProtocol {
pub id: String,
pub local: u16,
pub observed: u16,
}
impl ProtocolVersions {
pub fn snapshot(observed: &ObservedVersions) -> Self {
let local = local_manifest().protocols;
let observed_versions = lock(&observed.versions).clone();
let newer = observed_versions
.iter()
.filter_map(|(id, remote)| {
let local_version = local.get(id)?;
(*remote > *local_version).then(|| NewerProtocol {
id: id.clone(),
local: *local_version,
observed: *remote,
})
})
.collect();
Self {
local,
observed: observed_versions,
observed_peers: lock(&observed.peers).len(),
newer,
}
}
}
#[derive(Default)]
pub struct ObservedVersions {
versions: Mutex<BTreeMap<String, u16>>,
peers: Mutex<BTreeMap<String, String>>,
}
fn local_manifest() -> CapabilityManifest {
CapabilityManifest::frid("furumi", env!("CARGO_PKG_VERSION"))
.with_protocol("audio", super::audio::AUDIO_PROTOCOL_VERSION)
}
pub async fn serve(mut acceptor: StreamAcceptor) {
while let Some(stream) = acceptor.accept().await {
tokio::spawn(async move {
if let Err(error) = serve_one(stream).await {
tracing::debug!("capability stream failed: {error:#}");
}
});
}
}
async fn serve_one(mut stream: ByteStream) -> Result<()> {
let response = match read_message(&mut stream).await? {
CapabilityMessage::Get {
version: CAPABILITIES_PROTOCOL_VERSION,
} => CapabilityMessage::Manifest {
manifest: local_manifest(),
},
CapabilityMessage::Get { version } => CapabilityMessage::Error {
message: format!("unsupported capability protocol {version}"),
},
_ => CapabilityMessage::Error {
message: "expected capability request".to_string(),
},
};
write_message(&mut stream, &response).await?;
stream.send.finish()?;
let _ = tokio::time::timeout(Duration::from_secs(2), stream.send.stopped()).await;
Ok(())
}
pub async fn probe_loop(service: Arc<MusicDhtService>, observed: Arc<ObservedVersions>) {
let mut interval = tokio::time::interval(PROBE_INTERVAL);
loop {
interval.tick().await;
let peers = service
.connected_peers()
.into_iter()
.chain(
service
.known_peers()
.into_iter()
.map(|contact| contact.peer_id),
)
.collect::<std::collections::BTreeSet<_>>();
for peer in peers {
if let Err(error) = probe_peer(&service, peer, &observed).await {
tracing::trace!(%peer, "peer capability probe unavailable: {error:#}");
}
}
}
}
async fn probe_peer(
service: &MusicDhtService,
peer: EndpointId,
observed: &ObservedVersions,
) -> Result<()> {
let mut stream = service.open_stream(peer, CAPABILITIES_ALPN).await?;
write_message(
&mut stream,
&CapabilityMessage::Get {
version: CAPABILITIES_PROTOCOL_VERSION,
},
)
.await?;
stream.send.finish()?;
let response = tokio::time::timeout(Duration::from_secs(5), read_message(&mut stream))
.await
.context("capability request timed out")??;
let CapabilityMessage::Manifest { manifest } = response else {
anyhow::bail!("peer did not return a capability manifest");
};
manifest.validate()?;
{
let mut versions = lock(&observed.versions);
for (id, version) in manifest.protocols {
versions
.entry(id)
.and_modify(|current| *current = (*current).max(version))
.or_insert(version);
}
}
lock(&observed.peers).insert(peer.to_string(), manifest.application_version);
Ok(())
}
fn lock<T>(mutex: &Mutex<T>) -> MutexGuard<'_, T> {
mutex
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn local_manifest_lists_every_player_protocol() {
let manifest = local_manifest();
for id in [
"federation_net",
"ticket",
"rendezvous",
"music_dht",
"catalog",
"audio",
"device_sync",
"jam",
] {
assert!(manifest.protocols.contains_key(id), "missing {id}");
}
manifest.validate().unwrap();
}
#[test]
fn snapshot_reports_only_strictly_newer_versions() {
let observed = ObservedVersions::default();
lock(&observed.versions).insert("music_dht".to_string(), 99);
lock(&observed.versions).insert("jam".to_string(), crate::jam::PROTOCOL_VERSION);
let snapshot = ProtocolVersions::snapshot(&observed);
assert_eq!(snapshot.newer.len(), 1);
assert_eq!(snapshot.newer[0].id, "music_dht");
}
}
+136 -50
View File
@@ -13,6 +13,7 @@
//! the network too).
mod audio;
mod capabilities;
pub mod catalog;
use std::collections::{HashMap, VecDeque};
@@ -36,6 +37,7 @@ use crate::library::NetworkArtistPreview;
use crate::library::models::{ArtistRef, TrackItem};
pub use audio::{AUDIO_ALPN, DownloadProgress, StreamingStart, TrackMetadata};
pub use capabilities::ProtocolVersions;
pub use catalog::{CATALOG_ALPN, FedAppearsOn, FedArtistCard, FedCardTrack, FedRelease};
/// How often the published library is re-synchronized with the local index.
@@ -386,6 +388,7 @@ pub struct FedStatus {
pub last_sync: Option<String>,
pub last_error: Option<String>,
pub transport: TransportStatsSnapshot,
pub protocols: ProtocolVersions,
}
/// Outcome of preparing a federated track for playback.
@@ -419,13 +422,16 @@ pub struct Federation {
jam: Arc<crate::jam::JamManager>,
data_dir: PathBuf,
cache_dir: PathBuf,
media_dir: PathBuf,
media_dir: std::sync::Mutex<PathBuf>,
/// Serializes permanent downloads/card writes with a directory change.
media_change: tokio::sync::Mutex<()>,
metadata_cache: std::sync::Mutex<std::collections::HashMap<String, CachedTrackMetadata>>,
settings: std::sync::Mutex<FedSettings>,
running: tokio::sync::Mutex<Option<Running>>,
last_sync: std::sync::Mutex<Option<String>>,
last_error: std::sync::Mutex<Option<String>>,
transport_stats: Arc<TransportStats>,
observed_protocols: Arc<capabilities::ObservedVersions>,
}
#[derive(Debug, Clone)]
@@ -520,6 +526,7 @@ impl Federation {
library: Arc<Library>,
devices: Arc<crate::devices::DeviceSync>,
jam: Arc<crate::jam::JamManager>,
media_dir: PathBuf,
) -> Arc<Self> {
let dirs = crate::config::project_dirs();
let data_dir = dirs
@@ -530,10 +537,6 @@ impl Federation {
.as_ref()
.map(|d| d.cache_dir().join("fedcache"))
.unwrap_or_else(|| PathBuf::from("fedcache"));
let media_dir = dirs
.as_ref()
.map(|d| d.data_dir().join("federation-media"))
.unwrap_or_else(|| PathBuf::from("federation-media"));
let initial_error = [&data_dir, &cache_dir, &media_dir]
.into_iter()
.find_map(|dir| {
@@ -541,19 +544,22 @@ impl Federation {
.err()
.map(|err| format!("cannot create {}: {err}", dir.display()))
});
let media_dir = std::fs::canonicalize(&media_dir).unwrap_or(media_dir);
Arc::new(Self {
library,
devices,
jam,
data_dir,
cache_dir,
media_dir,
media_dir: std::sync::Mutex::new(media_dir),
media_change: tokio::sync::Mutex::new(()),
metadata_cache: std::sync::Mutex::new(Default::default()),
settings: std::sync::Mutex::new(load_settings()),
running: tokio::sync::Mutex::new(None),
last_sync: std::sync::Mutex::new(None),
last_error: std::sync::Mutex::new(initial_error),
transport_stats: Arc::new(TransportStats::default()),
observed_protocols: Arc::new(capabilities::ObservedVersions::default()),
})
}
@@ -561,6 +567,40 @@ impl Federation {
lock(&self.settings).clone()
}
pub fn media_dir(&self) -> PathBuf {
lock(&self.media_dir).clone()
}
/// Switches the permanent-download root while excluding concurrent
/// downloads. Validation runs again here because permissions may have
/// changed after the confirmation popup was shown.
pub async fn change_media_dir(
self: &Arc<Self>,
requested: PathBuf,
move_existing: bool,
) -> Result<crate::library::MusicRelocationStats> {
let _guard = self.media_change.lock().await;
let old = self.media_dir();
let library = Arc::clone(&self.library);
let requested_for_task = requested.clone();
let (path, stats) = tokio::task::spawn_blocking(move || -> Result<_> {
let path = Library::validate_music_directory(&requested_for_task)?;
if path == std::fs::canonicalize(&old).unwrap_or(old.clone()) {
return Ok((path, crate::library::MusicRelocationStats::default()));
}
let stats = if move_existing && old.exists() {
library.relocate_managed_music(&old, &path)?
} else {
crate::library::MusicRelocationStats::default()
};
Ok((path, stats))
})
.await
.context("music directory task failed")??;
*lock(&self.media_dir) = path;
Ok(stats)
}
pub async fn create_jam(&self) -> Result<String> {
let service = {
let running = self.running.lock().await;
@@ -656,6 +696,8 @@ impl Federation {
.stream_protocol(crate::devices::SYNC_ALPN)
// Capability-scoped shared playback control.
.stream_protocol(crate::jam::JAM_ALPN)
// Informational application/protocol versions.
.schema_independent_stream_protocol(capabilities::CAPABILITIES_ALPN)
.build()
.map_err(|err| anyhow::anyhow!("invalid federation config: {err}"))?;
let (service, mut events) = MusicDhtService::start(config)
@@ -728,6 +770,14 @@ impl Federation {
Arc::clone(&self.jam),
Arc::clone(&service),
));
let capabilities_acceptor = service
.stream_acceptor(capabilities::CAPABILITIES_ALPN)
.map_err(|err| anyhow::anyhow!("failed to take capabilities acceptor: {err}"))?;
let capabilities_serve_task = tokio::spawn(capabilities::serve(capabilities_acceptor));
let capabilities_probe_task = tokio::spawn(capabilities::probe_loop(
Arc::clone(&service),
Arc::clone(&self.observed_protocols),
));
*guard = Some(Running {
service,
@@ -742,6 +792,8 @@ impl Federation {
device_tick_task,
jam_serve_task,
jam_poll_task,
capabilities_serve_task,
capabilities_probe_task,
],
});
self.set_error(None);
@@ -880,6 +932,7 @@ impl Federation {
network: settings.network_id,
last_sync: lock(&self.last_sync).clone(),
last_error: lock(&self.last_error).clone(),
protocols: ProtocolVersions::snapshot(&self.observed_protocols),
..FedStatus::default()
};
if let Some(running) = guard.as_ref() {
@@ -1385,29 +1438,21 @@ impl Federation {
Ok(peer.to_string())
}
/// Directory for streamed (never library-imported) card artwork.
fn art_cache_dir(&self) -> PathBuf {
self.cache_dir.join("art")
}
/// Returns a cached-or-streamed image for the card: the artist image
/// (`release: None`) or a release cover. Peers are tried in order until
/// one answers with an image; the result lands in the art cache and its
/// local path is returned.
/// one answers with an image; the result is stored beside the artist or
/// release in the permanent music tree and its local path is returned.
pub async fn card_image(
&self,
owners: &[String],
artist: &str,
release: Option<&str>,
) -> Option<String> {
let dir = self.art_cache_dir();
let _guard = self.media_change.lock().await;
let dir = music_art_dir(&self.media_dir(), artist, release);
let stem = match release {
Some(release) => format!(
"cover-{}-{}",
sanitize_file_stem(artist),
sanitize_file_stem(release)
),
None => format!("artist-{}", sanitize_file_stem(artist)),
Some(_) => "cover".to_string(),
None => "artist".to_string(),
};
// Reuse a previously streamed copy of any known image type.
for extension in ["jpg", "png", "webp", "gif", "bmp"] {
@@ -1471,7 +1516,7 @@ impl Federation {
F: FnMut(DownloadProgress) + Send,
{
let save = self.settings().save_on_listen;
self.fetch_playable_with_progress(fed, save, false, progress, None)
self.fetch_playable_with_progress(fed, save, save, progress, None)
.await
}
@@ -1486,7 +1531,7 @@ impl Federation {
S: FnMut(StreamingStart) + Send,
{
let save = self.settings().save_on_listen;
self.fetch_playable_with_progress(fed, save, false, progress, Some(&mut stream_start))
self.fetch_playable_with_progress(fed, save, save, progress, Some(&mut stream_start))
.await
}
@@ -1540,6 +1585,11 @@ impl Federation {
where
F: FnMut(DownloadProgress) + Send,
{
let _media_guard = if save {
Some(self.media_change.lock().await)
} else {
None
};
let mut fed = fed.clone();
let mut tried_content_lookup = false;
@@ -1618,8 +1668,10 @@ impl Federation {
}
Err(_) => anyhow::bail!("malformed owner id '{}'", fed.owner),
};
let media_dir;
let dir = if save {
&self.media_dir
media_dir = music_release_dir(&self.media_dir(), &fed);
&media_dir
} else {
&self.cache_dir
};
@@ -1680,8 +1732,14 @@ impl Federation {
}
// Same for the cover: the peer's library cover wins over an
// embedded picture; embedded art stays as the fallback.
if import_cover.is_some() {
import.cover = import_cover;
if let Some((bytes, extension)) = &import_cover {
match save_release_cover(&import_path, bytes, extension) {
Ok(()) => import.cover = None,
Err(err) => {
tracing::warn!(%err, "saving the release cover beside its music failed");
import.cover = import_cover;
}
}
}
let (track_id, _) = crate::library::import::upsert_track(&library, &import)?;
// A like that referenced the federated track moves onto the
@@ -1693,7 +1751,13 @@ impl Federation {
// created (or still image-less) main artist.
if let (Some((bytes, extension)), Some(artist_name)) =
(&artist_image, import.artists.first())
&& let Err(err) = save_artist_image(&library, artist_name, bytes, extension)
&& let Err(err) = save_artist_image(
&library,
artist_name,
import_path.parent().and_then(Path::parent),
bytes,
extension,
)
{
tracing::warn!(%err, "saving the artist image failed");
}
@@ -1920,20 +1984,18 @@ impl Federation {
}
}
/// Writes a received artist image into the covers directory and attaches it
/// to the artist unless one is already set.
/// Writes a received artist image into the artist's music directory and
/// attaches it to the artist unless one is already set.
fn save_artist_image(
library: &Library,
artist_name: &str,
artist_dir: Option<&Path>,
bytes: &[u8],
extension: &str,
) -> Result<()> {
let covers_dir = library.covers_dir();
std::fs::create_dir_all(covers_dir)?;
let path = covers_dir.join(format!(
"artist-{}.{extension}",
sanitize_file_stem(artist_name)
));
let artist_dir = artist_dir.context("downloaded track has no artist directory")?;
std::fs::create_dir_all(artist_dir)?;
let path = artist_dir.join(format!("artist.{extension}"));
// Write only if the artist actually lacks an image, to avoid litter.
if library.artist_image_missing(artist_name)? {
std::fs::write(&path, bytes)?;
@@ -1942,6 +2004,16 @@ fn save_artist_image(
Ok(())
}
fn save_release_cover(audio_path: &Path, bytes: &[u8], extension: &str) -> Result<()> {
let directory = audio_path
.parent()
.context("downloaded track has no release directory")?;
std::fs::create_dir_all(directory)?;
let path = directory.join(format!("cover.{extension}"));
std::fs::write(&path, bytes).with_context(|| format!("writing {}", path.display()))?;
Ok(())
}
/// Overlays the peer-supplied metadata onto tag-derived import data. Every
/// non-empty peer field wins; file tags only fill the gaps.
fn apply_remote_metadata(import: &mut crate::library::import::TrackImport, meta: &TrackMetadata) {
@@ -2528,29 +2600,43 @@ fn collect_specs(library: &Library) -> Result<Vec<ItemSpec>> {
Ok(specs)
}
fn sanitize_file_stem(value: &str) -> String {
let cleaned: String = value
.chars()
.map(|c| match c {
'/' | '\\' | ':' | '*' | '?' | '"' | '<' | '>' | '|' => '_',
c if c.is_control() => '_',
c => c,
})
.collect();
let trimmed = cleaned.trim().trim_matches('.');
let mut stem: String = trimmed.chars().take(120).collect();
if stem.is_empty() {
stem.push_str("track");
fn music_artist_dir(root: &Path, artist: &str) -> PathBuf {
let artist = if artist.trim().is_empty() {
"Unknown Artist"
} else {
artist
};
root.join(crate::library::storage_name(artist, "Unknown Artist"))
}
fn music_art_dir(root: &Path, artist: &str, release: Option<&str>) -> PathBuf {
let artist_dir = music_artist_dir(root, artist);
match release.filter(|release| !release.trim().is_empty()) {
Some(release) => artist_dir.join(crate::library::storage_name(release, "Unknown Release")),
None => artist_dir,
}
stem
}
fn music_release_dir(root: &Path, fed: &FedTrack) -> PathBuf {
let artist = fed
.artist_names
.first()
.map(String::as_str)
.unwrap_or("Unknown Artist");
let release = fed
.release_title
.as_deref()
.filter(|release| !release.trim().is_empty())
.unwrap_or("Unknown Release");
music_art_dir(root, artist, Some(release))
}
fn download_stem(fed: &FedTrack) -> String {
let artists = fed.artist_line();
if artists.is_empty() {
sanitize_file_stem(&fed.title)
crate::library::storage_name(&fed.title, "track")
} else {
sanitize_file_stem(&format!("{artists} - {}", fed.title))
crate::library::storage_name(&format!("{artists} - {}", fed.title), "track")
}
}
+1 -5
View File
@@ -13,7 +13,7 @@ use crate::app::event::AppEvent;
use crate::devices::{PlaybackCommand, PlaybackSnapshot};
pub const JAM_ALPN: &[u8] = b"furumi/jam/1";
const PROTOCOL_VERSION: u16 = 1;
pub const PROTOCOL_VERSION: u16 = 1;
const MAX_LINE: usize = 8 * 1024 * 1024;
const MAX_COMMANDS: usize = 128;
const PARTICIPANT_TTL_MS: i64 = 30 * 60 * 1_000;
@@ -29,7 +29,6 @@ pub enum JamRole {
#[derive(Debug, Clone)]
pub struct JamStatus {
pub role: JamRole,
pub jam_id: Option<String>,
pub host_name: Option<String>,
pub invite: Option<String>,
pub participants: Vec<JamParticipant>,
@@ -41,7 +40,6 @@ impl Default for JamStatus {
fn default() -> Self {
Self {
role: JamRole::None,
jam_id: None,
host_name: None,
invite: None,
participants: Vec::new(),
@@ -207,7 +205,6 @@ impl JamManager {
if let Some(joined) = &state.joined {
return JamStatus {
role: JamRole::Participant,
jam_id: Some(joined.invite.jam_id.clone()),
host_name: Some(joined.invite.host_name.clone()),
invite: None,
participants: joined.participants.clone(),
@@ -218,7 +215,6 @@ impl JamManager {
if let Some(invite) = &state.host.invite {
return JamStatus {
role: JamRole::Host,
jam_id: Some(invite.jam_id.clone()),
host_name: Some(invite.host_name.clone()),
invite: state.host.invite_uri.clone(),
participants: state.host.participants.values().cloned().collect(),
+352
View File
@@ -11,6 +11,7 @@ pub mod import;
pub mod models;
use std::collections::{HashMap, HashSet};
use std::io::Write as _;
use std::path::{Path, PathBuf};
use std::sync::Mutex;
@@ -288,6 +289,12 @@ pub struct LocalLibraryStats {
pub database_bytes: u64,
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
pub struct MusicRelocationStats {
pub tracks: usize,
pub images: usize,
}
impl LocalLibraryStats {
pub fn total_bytes(&self) -> u64 {
self.audio_bytes
@@ -330,6 +337,261 @@ impl Library {
&self.covers_dir
}
/// Verifies that `path` can actually be used for durable downloads. The
/// returned path is absolute/canonical, so the persisted setting does not
/// later depend on Furumi's working directory.
pub fn validate_music_directory(path: &Path) -> Result<PathBuf> {
anyhow::ensure!(!path.as_os_str().is_empty(), "music directory is empty");
std::fs::create_dir_all(path)
.with_context(|| format!("creating music directory {}", path.display()))?;
let path = std::fs::canonicalize(path)
.with_context(|| format!("resolving music directory {}", path.display()))?;
anyhow::ensure!(path.is_dir(), "{} is not a directory", path.display());
let unique = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|duration| duration.as_nanos())
.unwrap_or(0);
let probe = path.join(format!(
".furumi-write-test-{}-{unique}",
std::process::id()
));
let result = (|| -> Result<()> {
let mut file = std::fs::OpenOptions::new()
.write(true)
.create_new(true)
.open(&probe)
.with_context(|| format!("no write access to {}", path.display()))?;
file.write_all(b"furumi")?;
file.sync_all()?;
Ok(())
})();
let _ = std::fs::remove_file(&probe);
result?;
Ok(path)
}
/// Copies Furumi-managed audio and its managed artwork into a portable
/// `Artist/Release` tree, atomically switches SQLite paths, then removes
/// the old copies. User-imported audio outside `old_root` is untouched.
pub fn relocate_managed_music(
&self,
old_root: &Path,
new_root: &Path,
) -> Result<MusicRelocationStats> {
let new_root = Self::validate_music_directory(new_root)?;
let old_root = std::fs::canonicalize(old_root)
.with_context(|| format!("resolving old music directory {}", old_root.display()))?;
anyhow::ensure!(
old_root != new_root,
"the new music directory is the current directory"
);
anyhow::ensure!(
!old_root.starts_with(&new_root) && !new_root.starts_with(&old_root),
"choose a directory outside the current music directory"
);
#[derive(Debug)]
struct Row {
track_id: i64,
track_path: String,
release_id: i64,
release_title: String,
cover_path: Option<String>,
artist_id: Option<i64>,
artist_name: String,
image_path: Option<String>,
}
let rows = {
let conn = self.lock();
let mut statement = conn.prepare(
"SELECT t.id, t.file_path, r.id, r.title, r.cover_path,
a.id, COALESCE(a.name, 'Unknown Artist'), a.image_path
FROM tracks t
JOIN releases r ON r.id = t.release_id
LEFT JOIN track_artists ta
ON ta.track_id = t.id AND ta.role = 'main' AND ta.position = 0
LEFT JOIN artists a ON a.id = ta.artist_id
ORDER BY t.id",
)?;
statement
.query_map([], |row| {
Ok(Row {
track_id: row.get(0)?,
track_path: row.get(1)?,
release_id: row.get(2)?,
release_title: row.get(3)?,
cover_path: row.get(4)?,
artist_id: row.get(5)?,
artist_name: row.get(6)?,
image_path: row.get(7)?,
})
})?
.collect::<rusqlite::Result<Vec<_>>>()?
};
#[derive(Debug, Clone)]
struct CopyOp {
source: PathBuf,
destination: PathBuf,
}
let managed_art =
|path: &Path| path.starts_with(&old_root) || path.starts_with(&self.covers_dir);
let mut reserved = HashSet::<PathBuf>::new();
let mut copies = Vec::<CopyOp>::new();
let mut track_updates = Vec::<(i64, String, String)>::new();
let mut release_updates = HashMap::<i64, (String, String)>::new();
let mut artist_updates = HashMap::<i64, (String, String)>::new();
for row in rows {
let source = PathBuf::from(&row.track_path);
let source = std::fs::canonicalize(&source).unwrap_or(source);
if !source.is_file() || !source.starts_with(&old_root) {
continue;
}
let artist_dir = new_root.join(storage_name(&row.artist_name, "Unknown Artist"));
let release_dir = artist_dir.join(storage_name(&row.release_title, "Unknown Release"));
let filename = source
.file_name()
.filter(|name| !name.is_empty())
.unwrap_or_else(|| std::ffi::OsStr::new("track"));
let destination =
unique_destination(release_dir.join(filename), row.track_id, &mut reserved);
copies.push(CopyOp {
source: source.clone(),
destination: destination.clone(),
});
track_updates.push((
row.track_id,
row.track_path,
destination.to_string_lossy().into_owned(),
));
if let Some(cover) = row.cover_path {
let cover_source = PathBuf::from(&cover);
let cover_source = std::fs::canonicalize(&cover_source).unwrap_or(cover_source);
if cover_source.is_file()
&& managed_art(&cover_source)
&& !release_updates.contains_key(&row.release_id)
{
let extension = cover_source.extension().unwrap_or_default();
let mut destination = release_dir.join("cover");
destination.set_extension(extension);
let destination =
unique_destination(destination, row.release_id, &mut reserved);
copies.push(CopyOp {
source: cover_source,
destination: destination.clone(),
});
release_updates.insert(
row.release_id,
(cover, destination.to_string_lossy().into_owned()),
);
}
}
if let (Some(artist_id), Some(image)) = (row.artist_id, row.image_path) {
let image_source = PathBuf::from(&image);
let image_source = std::fs::canonicalize(&image_source).unwrap_or(image_source);
if image_source.is_file()
&& managed_art(&image_source)
&& !artist_updates.contains_key(&artist_id)
{
let extension = image_source.extension().unwrap_or_default();
let mut destination = artist_dir.join("artist");
destination.set_extension(extension);
let destination = unique_destination(destination, artist_id, &mut reserved);
copies.push(CopyOp {
source: image_source,
destination: destination.clone(),
});
artist_updates.insert(
artist_id,
(image, destination.to_string_lossy().into_owned()),
);
}
}
}
let mut created = Vec::<PathBuf>::new();
let copy_result = (|| -> Result<()> {
for op in &copies {
let parent = op
.destination
.parent()
.context("music destination has no parent")?;
std::fs::create_dir_all(parent)?;
copy_file_exclusive(&op.source, &op.destination)?;
created.push(op.destination.clone());
}
Ok(())
})();
if let Err(error) = copy_result {
for path in created.iter().rev() {
let _ = std::fs::remove_file(path);
}
return Err(error.context("copying the existing music library"));
}
let database_result = (|| -> Result<()> {
let mut conn = self.lock();
let transaction = conn.transaction()?;
for (id, old, new) in &track_updates {
anyhow::ensure!(
transaction.execute(
"UPDATE tracks SET file_path = ?3 WHERE id = ?1 AND file_path = ?2",
params![id, old, new],
)? == 1,
"track {id} changed while the music library was moving"
);
}
for (id, (old, new)) in &release_updates {
anyhow::ensure!(
transaction.execute(
"UPDATE releases SET cover_path = ?3 WHERE id = ?1 AND cover_path = ?2",
params![id, old, new],
)? == 1,
"release {id} changed while the music library was moving"
);
}
for (id, (old, new)) in &artist_updates {
anyhow::ensure!(
transaction.execute(
"UPDATE artists SET image_path = ?3 WHERE id = ?1 AND image_path = ?2",
params![id, old, new],
)? == 1,
"artist {id} changed while the music library was moving"
);
}
transaction.commit()?;
Ok(())
})();
if let Err(error) = database_result {
for path in created.iter().rev() {
let _ = std::fs::remove_file(path);
}
return Err(error.context("updating music paths in the library"));
}
let mut removed = HashSet::new();
for op in &copies {
if removed.insert(op.source.clone())
&& let Err(error) = std::fs::remove_file(&op.source)
{
// The committed destination is authoritative. A failed old
// delete only leaves a recoverable duplicate.
tracing::warn!(path = %op.source.display(), %error, "old music copy was not removed");
}
}
remove_empty_directories(&old_root);
Ok(MusicRelocationStats {
tracks: track_updates.len(),
images: release_updates.len() + artist_updates.len(),
})
}
pub fn local_stats(&self) -> Result<LocalLibraryStats> {
let (artist_count, release_count, track_count, audio_bytes, tracks_without_size) = {
let conn = self.lock();
@@ -2498,6 +2760,96 @@ impl Library {
}
}
pub(crate) fn storage_name(value: &str, fallback: &str) -> String {
let cleaned: String = value
.chars()
.map(|character| match character {
'/' | '\\' | ':' | '*' | '?' | '"' | '<' | '>' | '|' => '_',
character if character.is_control() => '_',
character => character,
})
.collect();
let cleaned = cleaned.trim().trim_matches('.');
let mut shortened: String = cleaned.chars().take(120).collect();
if shortened.is_empty() {
return fallback.to_string();
}
let windows_base = shortened
.split('.')
.next()
.unwrap_or_default()
.to_ascii_uppercase();
let windows_reserved = matches!(windows_base.as_str(), "CON" | "PRN" | "AUX" | "NUL")
|| windows_base
.strip_prefix("COM")
.or_else(|| windows_base.strip_prefix("LPT"))
.is_some_and(|number| number.len() == 1 && matches!(number.as_bytes()[0], b'1'..=b'9'));
if windows_reserved {
shortened.insert(0, '_');
}
shortened
}
fn unique_destination(requested: PathBuf, id: i64, reserved: &mut HashSet<PathBuf>) -> PathBuf {
if !requested.exists() && reserved.insert(requested.clone()) {
return requested;
}
let stem = requested
.file_stem()
.unwrap_or_else(|| std::ffi::OsStr::new("file"))
.to_string_lossy();
let extension = requested.extension().map(|value| value.to_os_string());
let mut candidate = requested.with_file_name(format!("{stem}-{id}"));
if let Some(extension) = extension {
candidate.set_extension(extension);
}
let mut suffix = 2usize;
while candidate.exists() || !reserved.insert(candidate.clone()) {
candidate = requested.with_file_name(format!("{stem}-{id}-{suffix}"));
if let Some(extension) = requested.extension() {
candidate.set_extension(extension);
}
suffix += 1;
}
candidate
}
fn copy_file_exclusive(source: &Path, destination: &Path) -> Result<()> {
let mut input =
std::fs::File::open(source).with_context(|| format!("opening {}", source.display()))?;
let mut output = std::fs::OpenOptions::new()
.write(true)
.create_new(true)
.open(destination)
.with_context(|| format!("creating {}", destination.display()))?;
let result = std::io::copy(&mut input, &mut output)
.with_context(|| format!("copying {}", source.display()))
.and_then(|_| {
output
.sync_all()
.with_context(|| format!("syncing {}", destination.display()))
});
if let Err(error) = result {
drop(output);
let _ = std::fs::remove_file(destination);
return Err(error);
}
Ok(())
}
fn remove_empty_directories(root: &Path) {
let Ok(entries) = std::fs::read_dir(root) else {
return;
};
for entry in entries.flatten() {
let path = entry.path();
if path.is_dir() {
remove_empty_directories(&path);
}
}
let _ = std::fs::remove_dir(root);
}
fn cleanup_empty_releases(tx: &rusqlite::Transaction) -> rusqlite::Result<()> {
tx.execute(
"DELETE FROM releases WHERE id NOT IN (SELECT release_id FROM tracks)",
+132
View File
@@ -60,6 +60,138 @@ fn artist_filters(hide_featured_only: bool) -> crate::config::settings::LibraryF
}
}
fn unique_test_dir(label: &str) -> std::path::PathBuf {
let unique = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos();
std::env::temp_dir().join(format!("furumi-{label}-{}-{unique}", std::process::id()))
}
#[test]
fn managed_music_relocation_builds_artist_release_tree_and_updates_paths() {
let root = unique_test_dir("music-relocation");
let old = root.join("old");
let new = root.join("new");
let covers = root.join("covers");
std::fs::create_dir_all(&old).unwrap();
std::fs::create_dir_all(&covers).unwrap();
let audio = old.join("legacy.flac");
let cover = covers.join("release.jpg");
let artist_image = covers.join("artist.png");
std::fs::write(&audio, b"audio").unwrap();
std::fs::write(&cover, b"cover").unwrap();
std::fs::write(&artist_image, b"artist").unwrap();
let conn = Connection::open_in_memory().unwrap();
conn.pragma_update(None, "foreign_keys", "ON").unwrap();
register_norm_function(&conn).unwrap();
conn.execute_batch(SCHEMA).unwrap();
let library = Library {
conn: Mutex::new(conn),
db_path: root.join("library.db"),
covers_dir: covers,
};
let track_id = import::upsert_track(
&library,
&import::TrackImport {
file_path: audio.to_string_lossy().into_owned(),
title: "Song".into(),
artists: vec!["Artist".into()],
featured_artists: vec![],
album_artists: vec!["Artist".into()],
release_title: "Release".into(),
release_type: Some("album".into()),
year: Some(2026),
track_number: Some(1),
disc_number: Some(1),
duration_seconds: 1.0,
audio_format: Some("flac".into()),
audio_bitrate: None,
audio_sample_rate: None,
audio_bit_depth: None,
file_size_bytes: Some(5),
cover: None,
},
)
.unwrap()
.0;
let (release_id, artist_id): (i64, i64) = library
.lock()
.query_row(
"SELECT t.release_id, ta.artist_id
FROM tracks t JOIN track_artists ta ON ta.track_id = t.id
WHERE t.id = ?1 AND ta.role = 'main'",
[track_id],
|row| Ok((row.get(0)?, row.get(1)?)),
)
.unwrap();
library
.lock()
.execute(
"UPDATE releases SET cover_path = ?2 WHERE id = ?1",
params![release_id, cover.to_string_lossy()],
)
.unwrap();
library
.lock()
.execute(
"UPDATE artists SET image_path = ?2 WHERE id = ?1",
params![artist_id, artist_image.to_string_lossy()],
)
.unwrap();
let stats = library.relocate_managed_music(&old, &new).unwrap();
assert_eq!(stats.tracks, 1);
assert_eq!(stats.images, 2);
let track = library.tracks_by_ids(&[track_id]).unwrap().remove(0);
assert_eq!(
track.file_path,
new.join("Artist/Release/legacy.flac").to_string_lossy()
);
assert_eq!(
track.cover_path.as_deref(),
Some(
new.join("Artist/Release/cover.jpg")
.to_string_lossy()
.as_ref()
)
);
let image: String = library
.lock()
.query_row(
"SELECT image_path FROM artists WHERE id = ?1",
[artist_id],
|row| row.get(0),
)
.unwrap();
assert_eq!(image, new.join("Artist/artist.png").to_string_lossy());
assert!(!audio.exists());
std::fs::remove_dir_all(root).unwrap();
}
#[test]
fn music_directory_validation_rejects_a_file_without_touching_it() {
let root = unique_test_dir("music-validation");
std::fs::create_dir_all(&root).unwrap();
let file = root.join("not-a-directory");
std::fs::write(&file, b"keep").unwrap();
assert!(Library::validate_music_directory(&file).is_err());
assert_eq!(std::fs::read(&file).unwrap(), b"keep");
std::fs::remove_dir_all(root).unwrap();
}
#[test]
fn managed_music_names_are_portable_to_windows() {
assert_eq!(storage_name("Artist/Name", "fallback"), "Artist_Name");
assert_eq!(storage_name("CON", "fallback"), "_CON");
assert_eq!(storage_name("lpt9.live", "fallback"), "_lpt9.live");
assert_eq!(storage_name("...", "fallback"), "fallback");
}
#[test]
fn local_stats_counts_library_rows_and_audio_bytes() {
let lib = test_library();
+52 -1
View File
@@ -99,9 +99,10 @@ fn create_controls() -> Option<MediaControls> {
#[cfg(not(target_os = "windows"))]
let hwnd = None;
let dbus_name = platform_dbus_name();
let config = PlatformConfig {
display_name: "Furumi",
dbus_name: "cy.hexor.furumi",
dbus_name: &dbus_name,
hwnd,
};
match MediaControls::new(config) {
@@ -113,6 +114,41 @@ fn create_controls() -> Option<MediaControls> {
}
}
/// MPRIS well-known names must be unique on the session bus. Developers often
/// run a checkout next to an installed Furumi, and souvlaki 0.8.3 reports a
/// duplicate name by panicking in its private service thread. Ratatui's global
/// panic hook then restores the terminal even though the app thread is still
/// alive, leaving a frozen UI in canonical/echo mode.
///
/// Give every Unix MPRIS instance its own valid bus-name component. Other
/// backends ignore `dbus_name`, so retain the stable application identifier
/// there.
fn platform_dbus_name() -> String {
#[cfg(all(
unix,
not(any(target_os = "macos", target_os = "ios", target_os = "android"))
))]
{
mpris_dbus_name(std::process::id())
}
#[cfg(not(all(
unix,
not(any(target_os = "macos", target_os = "ios", target_os = "android"))
)))]
{
"cy.hexor.furumi".to_string()
}
}
#[cfg(all(
unix,
not(any(target_os = "macos", target_os = "ios", target_os = "android"))
))]
fn mpris_dbus_name(process_id: u32) -> String {
format!("cy.hexor.furumi.instance{process_id}")
}
/// An invisible top-level window owning the SMTC session. Created on the
/// main thread, which also pumps its messages in `pump_platform_events`.
#[cfg(target_os = "windows")]
@@ -217,3 +253,18 @@ fn pump_platform_events() {
#[cfg(not(any(target_os = "macos", target_os = "windows")))]
fn pump_platform_events() {}
#[cfg(test)]
mod tests {
#[cfg(all(
unix,
not(any(target_os = "macos", target_os = "ios", target_os = "android"))
))]
#[test]
fn mpris_names_are_distinct_between_processes() {
let first = super::mpris_dbus_name(41);
let second = super::mpris_dbus_name(42);
assert_eq!(first, "cy.hexor.furumi.instance41");
assert_ne!(first, second);
}
}
+184 -4
View File
@@ -2,6 +2,7 @@
use ratatui::Frame;
use ratatui::layout::{Constraint, Layout, Rect};
use ratatui::style::{Color, Modifier, Style};
use ratatui::text::{Line, Span};
use ratatui::widgets::{Block, Paragraph};
@@ -32,7 +33,7 @@ pub fn draw(frame: &mut Frame, area: Rect, state: &AppState) {
}
let rows_height =
(settings_rows(state).len() + 6 + device_presence_sections(state).len()) as u16;
(settings_rows(state).len() + 8 + device_presence_sections(state).len()) as u16;
let [rows_area, _, status_area] = Layout::vertical([
Constraint::Length(rows_height.min(inner.height)),
Constraint::Length(1),
@@ -68,6 +69,25 @@ fn draw_settings_rows(frame: &mut Frame, area: Rect, state: &AppState) {
let mut y = area.y;
let mut cursor = 0usize;
draw_section(frame, area, state, &mut y, "Library");
draw_row(
frame,
area,
state,
&mut y,
cursor,
state.settings_cursor,
"Music save directory",
if state.music_dir_changing {
format!("{} checking/changing…", state.spinner())
} else {
state.music_dir.to_string_lossy().into_owned()
},
);
cursor += 1;
y = y.saturating_add(1);
draw_section(frame, area, state, &mut y, "Federation");
for row in FedRow::ALL {
let (label, value) = match row {
@@ -338,6 +358,20 @@ fn draw_settings_rows(frame: &mut Frame, area: Rect, state: &AppState) {
);
}
fn protocol_label(id: &str) -> &str {
match id {
"federation_net" => "Federation transport",
"ticket" => "Peer ticket",
"rendezvous" => "Rendezvous",
"music_dht" => "Music DHT",
"catalog" => "Catalog",
"audio" => "Audio transfer",
"device_sync" => "Device sync",
"jam" => "Jam",
other => other,
}
}
fn draw_section(frame: &mut Frame, area: Rect, state: &AppState, y: &mut u16, title: &'static str) {
if *y >= area.y + area.height {
return;
@@ -485,11 +519,16 @@ fn draw_status(frame: &mut Frame, area: Rect, state: &AppState) {
return;
}
if area.width >= 60 && area.height >= 15 {
let [top_area, _, bottom_area] = Layout::vertical([
if area.width >= 60 && area.height >= 20 {
let protocols_height =
protocol_card_height(state, area.width.saturating_sub(2), area.height);
let [top_area, _, bottom_area, _, protocols_area, _] = Layout::vertical([
Constraint::Length(7),
Constraint::Length(1),
Constraint::Length(7),
Constraint::Length(1),
Constraint::Length(protocols_height),
Constraint::Min(0),
])
.areas(area);
let [status_area, _, local_area] = Layout::horizontal([
@@ -532,10 +571,17 @@ fn draw_status(frame: &mut Frame, area: Rect, state: &AppState) {
" Connected Devices ",
device_summary_lines(state),
);
draw_summary_card(
frame,
protocols_area,
state,
" Protocol Versions ",
protocol_summary_lines(state, protocols_area.width.saturating_sub(2)),
);
return;
}
if area.height < 31 {
if area.height < 39 {
frame.render_widget(
Paragraph::new(compact_status_lines(state))
.wrap(ratatui::widgets::Wrap { trim: false }),
@@ -553,6 +599,8 @@ fn draw_status(frame: &mut Frame, area: Rect, state: &AppState) {
_,
local_area,
_,
protocols_area,
_,
] = Layout::vertical([
Constraint::Length(7),
Constraint::Length(1),
@@ -561,6 +609,12 @@ fn draw_status(frame: &mut Frame, area: Rect, state: &AppState) {
Constraint::Length(7),
Constraint::Length(1),
Constraint::Length(7),
Constraint::Length(1),
Constraint::Length(protocol_card_height(
state,
area.width.saturating_sub(2),
area.height,
)),
Constraint::Min(0),
])
.areas(area);
@@ -593,6 +647,13 @@ fn draw_status(frame: &mut Frame, area: Rect, state: &AppState) {
" Local Data ",
local_data_summary_lines(state),
);
draw_summary_card(
frame,
protocols_area,
state,
" Protocol Versions ",
protocol_summary_lines(state, protocols_area.width.saturating_sub(2)),
);
}
fn compact_status_lines(state: &AppState) -> Vec<Line<'static>> {
@@ -608,6 +669,9 @@ fn compact_status_lines(state: &AppState) -> Vec<Line<'static>> {
lines.push(Line::default());
lines.push(Line::styled("Connected Devices", theme::header_for(state)));
lines.extend(device_summary_lines(state).into_iter().take(2));
lines.push(Line::default());
lines.push(Line::styled("Protocol Versions", theme::header_for(state)));
lines.extend(protocol_summary_lines(state, 0).into_iter().take(3));
lines
}
@@ -637,6 +701,122 @@ fn summary_line(label: &'static str, value: String) -> Line<'static> {
])
}
fn protocol_summary_line(
label: &str,
value: String,
style: Style,
label_width: usize,
) -> Line<'static> {
Line::from(vec![
Span::styled(
format!("{:<label_width$}", protocol_label(label)),
theme::dim(),
),
Span::styled(value, style),
])
}
fn protocol_summary_lines(state: &AppState, width: u16) -> Vec<Line<'static>> {
let Some(status) = state.federation.status.as_ref() else {
return vec![protocol_summary_line(
"status",
"[UNKNOWN]".to_string(),
theme::dim(),
22,
)];
};
let protocols = &status.protocols;
let newer = !protocols.newer.is_empty();
let badge = if newer {
"[NEWER VERSION SEEN]"
} else if status.running && protocols.observed_peers == 0 {
"[CURRENT · waiting for peers]"
} else {
"[CURRENT]"
};
let badge_style = Style::new()
.fg(if newer { Color::LightRed } else { Color::Green })
.add_modifier(Modifier::BOLD);
let mut lines = vec![protocol_summary_line(
"status",
badge.to_string(),
badge_style,
22,
)];
let mut entries = Vec::new();
for (id, local) in &protocols.local {
let observed = protocols.observed.get(id).copied();
let value = match observed {
Some(remote) if remote > *local => format!("local {local} · network {remote}"),
Some(remote) => format!("{local} · seen {remote}"),
None => local.to_string(),
};
let style = if observed.is_some_and(|remote| remote > *local) {
Style::new()
.fg(Color::LightRed)
.add_modifier(Modifier::BOLD)
} else {
Style::default()
};
entries.push((id.as_str(), value, style));
}
if width >= 58 {
let cell_width = width as usize / 2;
let label_width = 22.min(cell_width.saturating_sub(3));
let value_width = cell_width.saturating_sub(label_width + 2);
for pair in entries.chunks(2) {
let mut spans =
protocol_cell_spans(pair[0].0, &pair[0].1, pair[0].2, label_width, value_width);
if let Some(second) = pair.get(1) {
spans.push(Span::styled(" ", theme::dim()));
spans.extend(protocol_cell_spans(
second.0,
&second.1,
second.2,
label_width,
value_width,
));
}
lines.push(Line::from(spans));
}
} else {
lines.extend(
entries
.into_iter()
.map(|(id, value, style)| protocol_summary_line(id, value, style, 22)),
);
}
if newer {
lines.push(Line::styled(
"A newer protocol was observed; update Furumi for compatibility.",
Style::new().fg(Color::LightRed),
));
}
lines
}
fn protocol_card_height(state: &AppState, width: u16, available: u16) -> u16 {
let content = protocol_summary_lines(state, width).len() as u16;
content.saturating_add(2).min(available)
}
fn protocol_cell_spans(
label: &str,
value: &str,
style: Style,
label_width: usize,
value_width: usize,
) -> Vec<Span<'static>> {
let value = value.chars().take(value_width).collect::<String>();
vec![
Span::styled(
format!("{:<label_width$}", protocol_label(label)),
theme::dim(),
),
Span::styled(format!("{value:<value_width$}"), style),
]
}
fn node_summary_lines(state: &AppState) -> Vec<Line<'static>> {
match &state.federation.status {
None => vec![
+1 -1
View File
@@ -258,7 +258,7 @@ fn draw_queue(frame: &mut Frame, area: Rect, state: &AppState) {
let player = &state.player;
let block = Block::bordered()
.title(format!(
" Queue — {} tracks · Mode: {} · enter: play · d: remove · shift-v: select · :clear ",
" Queue — {} tracks · Mode: {} · enter: play · alt-j/k: move · d: remove · shift-v: select · :clear ",
player.queue.len(),
state.global.filters.source_mode.label()
))
+62 -9
View File
@@ -42,8 +42,15 @@ pub fn draw(frame: &mut Frame, state: &AppState) {
draw_fed_input(frame, state, field.title(), field.help(), input)
}
Some(Popup::FedText { title, text }) => draw_fed_text(frame, state, title, text),
Some(Popup::FedCopyText { title, text, help }) => {
draw_fed_copy_text(frame, state, title, text, help)
Some(Popup::FedCopyText {
title,
text,
help,
cursor,
}) => draw_fed_copy_text(frame, state, title, text, help, *cursor),
Some(Popup::PlainText { .. }) => {}
Some(Popup::ConfirmMusicDirectory { path }) => {
draw_music_directory_confirmation(frame, state, path)
}
Some(Popup::FederationStatusDetails {
focus,
@@ -93,6 +100,34 @@ pub fn draw(frame: &mut Frame, state: &AppState) {
}
}
fn draw_music_directory_confirmation(frame: &mut Frame, state: &AppState, path: &std::path::Path) {
let area = centered(frame.area(), 76, 9);
let block = Block::bordered()
.title(" Move existing music? ")
.title_style(theme::header_for(state))
.border_style(theme::strong_border_for(state));
let inner = block.inner(area);
frame.render_widget(Clear, area);
frame.render_widget(block, area);
frame.render_widget(
Paragraph::new(vec![
Line::raw("The destination is writable:"),
Line::styled(
path.to_string_lossy().into_owned(),
theme::accent_for(state),
),
Line::raw(""),
Line::raw("Move music and artwork previously saved by Furumi there?"),
Line::styled(
"enter/y move · n keep existing files where they are · esc cancel",
theme::dim(),
),
])
.wrap(Wrap { trim: false }),
inner,
);
}
fn draw_listen_history(frame: &mut Frame, state: &AppState, cursor: usize) {
let area = centered(
frame.area(),
@@ -862,7 +897,14 @@ fn draw_fed_text(frame: &mut Frame, state: &AppState, title: &str, text: &str) {
}
/// Wrapped text with an explicit copy-and-close action.
fn draw_fed_copy_text(frame: &mut Frame, state: &AppState, title: &str, text: &str, help: &str) {
fn draw_fed_copy_text(
frame: &mut Frame,
state: &AppState,
title: &str,
text: &str,
help: &str,
cursor: usize,
) {
let width = frame.area().width.saturating_sub(8).clamp(36, 96);
let text_width = usize::from(width.saturating_sub(2));
let text_lines = (text.chars().count() / text_width.max(1) + 1) as u16;
@@ -894,17 +936,28 @@ fn draw_fed_copy_text(frame: &mut Frame, state: &AppState, title: &str, text: &s
Paragraph::new(text.to_string()).wrap(Wrap { trim: false }),
body_area,
);
let button = |label: &str, selected: bool| {
if selected {
Span::styled(format!(" {label} "), theme::tab_active_for(state))
} else {
Span::styled(format!(" {label} "), theme::dim())
}
};
frame.render_widget(
Paragraph::new(Line::styled(
" Copy to clipboard and close ",
theme::tab_active_for(state),
))
Paragraph::new(Line::from(vec![
button("Copy to clipboard", cursor == 0),
Span::raw(" "),
button("Show as plain terminal line", cursor == 1),
]))
.alignment(Alignment::Center),
button_area,
);
frame.render_widget(
Paragraph::new(Line::styled("enter/c copy · esc close", theme::dim()))
.alignment(Alignment::Center),
Paragraph::new(Line::styled(
"left/right choose · enter apply · c copy · p plain line · esc close",
theme::dim(),
))
.alignment(Alignment::Center),
footer_area,
);
}
+23
View File
@@ -0,0 +1,23 @@
[package]
name = "block"
version = "0.1.6"
authors = ["Steven Sheldon"]
description = "Rust interface for Apple's C language extension of blocks."
keywords = ["blocks", "osx", "ios", "objective-c"]
readme = "README.md"
repository = "http://github.com/SSheldon/rust-block"
documentation = "http://ssheldon.github.io/rust-objc/block/"
license = "MIT"
exclude = [
".gitignore",
".travis.yml",
"travis_install.sh",
"travis_test.sh",
"tests-ios/**",
]
[dev-dependencies.objc_test_utils]
version = "0.0"
path = "test_utils"
+42
View File
@@ -0,0 +1,42 @@
Rust interface for Apple's C language extension of blocks.
For more information on the specifics of the block implementation, see
Clang's documentation: http://clang.llvm.org/docs/Block-ABI-Apple.html
## Invoking blocks
The `Block` struct is used for invoking blocks from Objective-C. For example,
consider this Objective-C function:
``` objc
int32_t sum(int32_t (^block)(int32_t, int32_t)) {
return block(5, 8);
}
```
We could write it in Rust as the following:
``` rust
unsafe fn sum(block: &Block<(i32, i32), i32>) -> i32 {
block.call((5, 8))
}
```
Note the extra parentheses in the `call` method, since the arguments must be
passed as a tuple.
## Creating blocks
Creating a block to pass to Objective-C can be done with the `ConcreteBlock`
struct. For example, to create a block that adds two `i32`s, we could write:
``` rust
let block = ConcreteBlock::new(|a: i32, b: i32| a + b);
let block = block.copy();
assert!(unsafe { block.call((5, 8)) } == 13);
```
It is important to copy your block to the heap (with the `copy` method) before
passing it to Objective-C; this is because our `ConcreteBlock` is only meant
to be copied once, and we can enforce this in Rust, but if Objective-C code
were to copy it twice we could have a double free.
+399
View File
@@ -0,0 +1,399 @@
/*!
A Rust interface for Objective-C blocks.
For more information on the specifics of the block implementation, see
Clang's documentation: http://clang.llvm.org/docs/Block-ABI-Apple.html
# Invoking blocks
The `Block` struct is used for invoking blocks from Objective-C. For example,
consider this Objective-C function:
``` objc
int32_t sum(int32_t (^block)(int32_t, int32_t)) {
return block(5, 8);
}
```
We could write it in Rust as the following:
```
# use block::Block;
unsafe fn sum(block: &Block<(i32, i32), i32>) -> i32 {
block.call((5, 8))
}
```
Note the extra parentheses in the `call` method, since the arguments must be
passed as a tuple.
# Creating blocks
Creating a block to pass to Objective-C can be done with the `ConcreteBlock`
struct. For example, to create a block that adds two `i32`s, we could write:
```
# use block::ConcreteBlock;
let block = ConcreteBlock::new(|a: i32, b: i32| a + b);
let block = block.copy();
assert!(unsafe { block.call((5, 8)) } == 13);
```
It is important to copy your block to the heap (with the `copy` method) before
passing it to Objective-C; this is because our `ConcreteBlock` is only meant
to be copied once, and we can enforce this in Rust, but if Objective-C code
were to copy it twice we could have a double free.
*/
#[cfg(test)]
mod test_utils;
use std::marker::PhantomData;
use std::mem;
use std::ops::{Deref, DerefMut};
use std::os::raw::{c_int, c_ulong, c_void};
use std::ptr;
#[repr(C)]
struct Class {
_private: [u8; 0],
}
#[cfg_attr(any(target_os = "macos", target_os = "ios"),
link(name = "System", kind = "dylib"))]
#[cfg_attr(not(any(target_os = "macos", target_os = "ios")),
link(name = "BlocksRuntime", kind = "dylib"))]
extern "C" {
static _NSConcreteStackBlock: Class;
fn _Block_copy(block: *const c_void) -> *mut c_void;
fn _Block_release(block: *const c_void);
}
/// Types that may be used as the arguments to an Objective-C block.
pub trait BlockArguments: Sized {
/// Calls the given `Block` with self as the arguments.
///
/// Unsafe because `block` must point to a valid `Block` and this invokes
/// foreign code whose safety the compiler cannot verify.
unsafe fn call_block<R>(self, block: *mut Block<Self, R>) -> R;
}
macro_rules! block_args_impl {
($($a:ident : $t:ident),*) => (
impl<$($t),*> BlockArguments for ($($t,)*) {
unsafe fn call_block<R>(self, block: *mut Block<Self, R>) -> R {
let invoke: unsafe extern "C" fn(*mut Block<Self, R> $(, $t)*) -> R = {
let base = block as *mut BlockBase<Self, R>;
mem::transmute((*base).invoke)
};
let ($($a,)*) = self;
invoke(block $(, $a)*)
}
}
);
}
block_args_impl!();
block_args_impl!(a: A);
block_args_impl!(a: A, b: B);
block_args_impl!(a: A, b: B, c: C);
block_args_impl!(a: A, b: B, c: C, d: D);
block_args_impl!(a: A, b: B, c: C, d: D, e: E);
block_args_impl!(a: A, b: B, c: C, d: D, e: E, f: F);
block_args_impl!(a: A, b: B, c: C, d: D, e: E, f: F, g: G);
block_args_impl!(a: A, b: B, c: C, d: D, e: E, f: F, g: G, h: H);
block_args_impl!(a: A, b: B, c: C, d: D, e: E, f: F, g: G, h: H, i: I);
block_args_impl!(a: A, b: B, c: C, d: D, e: E, f: F, g: G, h: H, i: I, j: J);
block_args_impl!(a: A, b: B, c: C, d: D, e: E, f: F, g: G, h: H, i: I, j: J, k: K);
block_args_impl!(a: A, b: B, c: C, d: D, e: E, f: F, g: G, h: H, i: I, j: J, k: K, l: L);
#[repr(C)]
struct BlockBase<A, R> {
isa: *const Class,
flags: c_int,
_reserved: c_int,
invoke: unsafe extern "C" fn(*mut Block<A, R>, ...) -> R,
}
/// An Objective-C block that takes arguments of `A` when called and
/// returns a value of `R`.
#[repr(C)]
pub struct Block<A, R> {
_base: PhantomData<BlockBase<A, R>>,
}
impl<A: BlockArguments, R> Block<A, R> where A: BlockArguments {
/// Call self with the given arguments.
///
/// Unsafe because this invokes foreign code that the caller must verify
/// doesn't violate any of Rust's safety rules. For example, if this block
/// is shared with multiple references, the caller must ensure that calling
/// it will not cause a data race.
pub unsafe fn call(&self, args: A) -> R {
args.call_block(self as *const _ as *mut _)
}
}
/// A reference-counted Objective-C block.
pub struct RcBlock<A, R> {
ptr: *mut Block<A, R>,
}
impl<A, R> RcBlock<A, R> {
/// Construct an `RcBlock` for the given block without copying it.
/// The caller must ensure the block has a +1 reference count.
///
/// Unsafe because `ptr` must point to a valid `Block` and must have a +1
/// reference count or it will be overreleased when the `RcBlock` is
/// dropped.
pub unsafe fn new(ptr: *mut Block<A, R>) -> Self {
RcBlock { ptr: ptr }
}
/// Constructs an `RcBlock` by copying the given block.
///
/// Unsafe because `ptr` must point to a valid `Block`.
pub unsafe fn copy(ptr: *mut Block<A, R>) -> Self {
let ptr = _Block_copy(ptr as *const c_void) as *mut Block<A, R>;
RcBlock { ptr: ptr }
}
}
impl<A, R> Clone for RcBlock<A, R> {
fn clone(&self) -> RcBlock<A, R> {
unsafe {
RcBlock::copy(self.ptr)
}
}
}
impl<A, R> Deref for RcBlock<A, R> {
type Target = Block<A, R>;
fn deref(&self) -> &Block<A, R> {
unsafe { &*self.ptr }
}
}
impl<A, R> Drop for RcBlock<A, R> {
fn drop(&mut self) {
unsafe {
_Block_release(self.ptr as *const c_void);
}
}
}
/// Types that may be converted into a `ConcreteBlock`.
pub trait IntoConcreteBlock<A>: Sized where A: BlockArguments {
/// The return type of the resulting `ConcreteBlock`.
type Ret;
/// Consumes self to create a `ConcreteBlock`.
fn into_concrete_block(self) -> ConcreteBlock<A, Self::Ret, Self>;
}
macro_rules! concrete_block_impl {
($f:ident) => (
concrete_block_impl!($f,);
);
($f:ident, $($a:ident : $t:ident),*) => (
impl<$($t,)* R, X> IntoConcreteBlock<($($t,)*)> for X
where X: Fn($($t,)*) -> R {
type Ret = R;
fn into_concrete_block(self) -> ConcreteBlock<($($t,)*), R, X> {
unsafe extern "C" fn $f<$($t,)* R, X>(
block_ptr: *mut ConcreteBlock<($($t,)*), R, X>
$(, $a: $t)*) -> R
where X: Fn($($t,)*) -> R {
let block = &*block_ptr;
(block.closure)($($a),*)
}
let f: unsafe extern "C" fn(*mut ConcreteBlock<($($t,)*), R, X> $(, $a: $t)*) -> R = $f;
unsafe {
ConcreteBlock::with_invoke(mem::transmute(f), self)
}
}
}
);
}
concrete_block_impl!(concrete_block_invoke_args0);
concrete_block_impl!(concrete_block_invoke_args1, a: A);
concrete_block_impl!(concrete_block_invoke_args2, a: A, b: B);
concrete_block_impl!(concrete_block_invoke_args3, a: A, b: B, c: C);
concrete_block_impl!(concrete_block_invoke_args4, a: A, b: B, c: C, d: D);
concrete_block_impl!(concrete_block_invoke_args5, a: A, b: B, c: C, d: D, e: E);
concrete_block_impl!(concrete_block_invoke_args6, a: A, b: B, c: C, d: D, e: E, f: F);
concrete_block_impl!(concrete_block_invoke_args7, a: A, b: B, c: C, d: D, e: E, f: F, g: G);
concrete_block_impl!(concrete_block_invoke_args8, a: A, b: B, c: C, d: D, e: E, f: F, g: G, h: H);
concrete_block_impl!(concrete_block_invoke_args9, a: A, b: B, c: C, d: D, e: E, f: F, g: G, h: H, i: I);
concrete_block_impl!(concrete_block_invoke_args10, a: A, b: B, c: C, d: D, e: E, f: F, g: G, h: H, i: I, j: J);
concrete_block_impl!(concrete_block_invoke_args11, a: A, b: B, c: C, d: D, e: E, f: F, g: G, h: H, i: I, j: J, k: K);
concrete_block_impl!(concrete_block_invoke_args12, a: A, b: B, c: C, d: D, e: E, f: F, g: G, h: H, i: I, j: J, k: K, l: L);
/// An Objective-C block whose size is known at compile time and may be
/// constructed on the stack.
#[repr(C)]
pub struct ConcreteBlock<A, R, F> {
base: BlockBase<A, R>,
descriptor: Box<BlockDescriptor<ConcreteBlock<A, R, F>>>,
closure: F,
}
impl<A, R, F> ConcreteBlock<A, R, F>
where A: BlockArguments, F: IntoConcreteBlock<A, Ret=R> {
/// Constructs a `ConcreteBlock` with the given closure.
/// When the block is called, it will return the value that results from
/// calling the closure.
pub fn new(closure: F) -> Self {
closure.into_concrete_block()
}
}
impl<A, R, F> ConcreteBlock<A, R, F> {
/// Constructs a `ConcreteBlock` with the given invoke function and closure.
/// Unsafe because the caller must ensure the invoke function takes the
/// correct arguments.
unsafe fn with_invoke(invoke: unsafe extern "C" fn(*mut Self, ...) -> R,
closure: F) -> Self {
ConcreteBlock {
base: BlockBase {
isa: &_NSConcreteStackBlock,
// 1 << 25 = BLOCK_HAS_COPY_DISPOSE
flags: 1 << 25,
_reserved: 0,
invoke: mem::transmute(invoke),
},
descriptor: Box::new(BlockDescriptor::new()),
closure: closure,
}
}
}
impl<A, R, F> ConcreteBlock<A, R, F> where F: 'static {
/// Copy self onto the heap as an `RcBlock`.
pub fn copy(self) -> RcBlock<A, R> {
unsafe {
let mut block = self;
let copied = RcBlock::copy(&mut *block);
// At this point, our copy helper has been run so the block will
// be moved to the heap and we can forget the original block
// because the heap block will drop in our dispose helper.
mem::forget(block);
copied
}
}
}
impl<A, R, F> Clone for ConcreteBlock<A, R, F> where F: Clone {
fn clone(&self) -> Self {
unsafe {
ConcreteBlock::with_invoke(mem::transmute(self.base.invoke),
self.closure.clone())
}
}
}
impl<A, R, F> Deref for ConcreteBlock<A, R, F> {
type Target = Block<A, R>;
fn deref(&self) -> &Block<A, R> {
unsafe { &*(&self.base as *const _ as *const Block<A, R>) }
}
}
impl<A, R, F> DerefMut for ConcreteBlock<A, R, F> {
fn deref_mut(&mut self) -> &mut Block<A, R> {
unsafe { &mut *(&mut self.base as *mut _ as *mut Block<A, R>) }
}
}
unsafe extern "C" fn block_context_dispose<B>(block: &mut B) {
// Read the block onto the stack and let it drop
ptr::read(block);
}
unsafe extern "C" fn block_context_copy<B>(_dst: &mut B, _src: &B) {
// The runtime memmoves the src block into the dst block, nothing to do
}
#[repr(C)]
struct BlockDescriptor<B> {
_reserved: c_ulong,
block_size: c_ulong,
copy_helper: unsafe extern "C" fn(&mut B, &B),
dispose_helper: unsafe extern "C" fn(&mut B),
}
impl<B> BlockDescriptor<B> {
fn new() -> BlockDescriptor<B> {
BlockDescriptor {
_reserved: 0,
block_size: mem::size_of::<B>() as c_ulong,
copy_helper: block_context_copy::<B>,
dispose_helper: block_context_dispose::<B>,
}
}
}
#[cfg(test)]
mod tests {
use test_utils::*;
use super::{ConcreteBlock, RcBlock};
#[test]
fn test_call_block() {
let block = get_int_block_with(13);
unsafe {
assert!(block.call(()) == 13);
}
}
#[test]
fn test_call_block_args() {
let block = get_add_block_with(13);
unsafe {
assert!(block.call((2,)) == 15);
}
}
#[test]
fn test_create_block() {
let block = ConcreteBlock::new(|| 13);
let result = invoke_int_block(&block);
assert!(result == 13);
}
#[test]
fn test_create_block_args() {
let block = ConcreteBlock::new(|a: i32| a + 5);
let result = invoke_add_block(&block, 6);
assert!(result == 11);
}
#[test]
fn test_concrete_block_copy() {
let s = "Hello!".to_string();
let expected_len = s.len() as i32;
let block = ConcreteBlock::new(move || s.len() as i32);
assert!(invoke_int_block(&block) == expected_len);
let copied = block.copy();
assert!(invoke_int_block(&copied) == expected_len);
}
#[test]
fn test_concrete_block_stack_copy() {
fn make_block() -> RcBlock<(), i32> {
let x = 7;
let block = ConcreteBlock::new(move || x);
block.copy()
}
let block = make_block();
assert!(invoke_int_block(&block) == 7);
}
}
+31
View File
@@ -0,0 +1,31 @@
extern crate objc_test_utils;
use {Block, RcBlock};
pub fn get_int_block_with(i: i32) -> RcBlock<(), i32> {
unsafe {
let ptr = objc_test_utils::get_int_block_with(i);
RcBlock::new(ptr as *mut _)
}
}
pub fn get_add_block_with(i: i32) -> RcBlock<(i32,), i32> {
unsafe {
let ptr = objc_test_utils::get_add_block_with(i);
RcBlock::new(ptr as *mut _)
}
}
pub fn invoke_int_block(block: &Block<(), i32>) -> i32 {
let ptr = block as *const _;
unsafe {
objc_test_utils::invoke_int_block(ptr as *mut _)
}
}
pub fn invoke_add_block(block: &Block<(i32,), i32>, a: i32) -> i32 {
let ptr = block as *const _;
unsafe {
objc_test_utils::invoke_add_block(ptr as *mut _, a)
}
}