Fix config loading, modularize crates, and add valid config template

This commit is contained in:
2026-07-21 21:09:04 +02:00
parent 66844ea01e
commit 4b731d406a
22 changed files with 6084 additions and 0 deletions
+5
View File
@@ -0,0 +1,5 @@
.cache/
.cache
config.toml
!config.example.toml
/target
Generated
+4236
View File
File diff suppressed because it is too large Load Diff
+63
View File
@@ -0,0 +1,63 @@
[workspace]
resolver = "2"
members = [
"crates/riptune-core",
"crates/riptune-cache",
"crates/riptune-mpris",
"crates/riptune-niri",
"crates/riptune-tui",
"crates/riptune-types",
]
[workspace.package]
version = "0.1.0"
edition = "2021"
license = "GPL-3.0"
repository = "https://github.com/Quinta0/riptune"
authors = ["Quintavalle Pietro"]
[workspace.dependencies]
tokio = { version = "1", features = ["full"] }
serde = { version = "1", features = ["derive"] }
serde_json = "1"
reqwest = { version = "0.12", default-features = false, features = ["json", "stream", "rustls-tls"] }
anyhow = "1"
thiserror = "1"
tracing = "0.1"
tracing-subscriber = "0.3"
[package]
name = "riptune"
version.workspace = true
edition.workspace = true
license.workspace = true
description = "A niri-native TUI music client for Subsonic servers. Rip and tune."
[[bin]]
name = "riptune"
path = "src/main.rs"
[dependencies]
riptune-core = { path = "crates/riptune-core" }
riptune-cache = { path = "crates/riptune-cache" }
riptune-mpris = { path = "crates/riptune-mpris" }
riptune-niri = { path = "crates/riptune-niri" }
riptune-tui = { path = "crates/riptune-tui" }
riptune-types = { path = "crates/riptune-types" }
tokio = { workspace = true }
anyhow = { workspace = true }
tracing = { workspace = true }
tracing-subscriber = { workspace = true }
serde = { workspace = true }
toml = "0.8"
clap = { version = "4", features = ["derive"] }
# Pinned exactly: rodio 0.21 reworked the OutputStream/Sink API around a new
# Mixer type (Sink::try_new -> Sink::connect_new, etc). Pinning avoids that
# rewrite landing under us; bump deliberately and update audio.rs to match.
rodio = { version = "=0.20.1", features = ["symphonia-all"] }
# "blocking" is additive on top of the workspace-defined reqwest (which has
# default-features = false + rustls-tls). Used only by the audio thread's
# synchronous fetch -- see the comment in src/audio.rs for why blocking is
# actually the right call there, not a shortcut.
reqwest = { workspace = true, features = ["blocking"] }
+11
View File
@@ -0,0 +1,11 @@
[server]
url = "https://your-navidrome-server.com"
username = "your_username"
password = "your_password_or_api_key"
[niri]
enabled = true
workspace_notifications = true
[cache]
path = "~/.cache/riptune/library.sqlite"
+13
View File
@@ -0,0 +1,13 @@
[package]
name = "riptune-cache"
version.workspace = true
edition.workspace = true
license.workspace = true
description = "SQLite-backed local index for instant TUI startup on large libraries"
[dependencies]
riptune-core = { path = "../riptune-core" }
tokio = { workspace = true }
thiserror = { workspace = true }
tracing = { workspace = true }
sqlx = { version = "0.7", features = ["runtime-tokio", "sqlite"] }
+244
View File
@@ -0,0 +1,244 @@
//! Local SQLite index of the remote Subsonic library.
//!
//! Why this exists: hitting `getArtists`/`getAlbum` over the network on every TUI
//! launch is slow on large libraries (multi-second waits are common on Airsonic
//! instances with 50k+ tracks). We index artists/albums/tracks locally and treat
//! the network as a background sync source, not the read path for rendering[cite: 13].
use riptune_core::SubsonicClient;
use riptune_core::models::{Artist, Track, AlbumSummary};
use sqlx::sqlite::SqlitePoolOptions;
use sqlx::SqlitePool;
use sqlx::Row;
use std::path::Path;
use thiserror::Error;
#[derive(Debug, Error)]
pub enum CacheError {
#[error("database error: {0}")]
Db(#[from] sqlx::Error),
#[error("io error: {0}")]
Io(#[from] std::io::Error),
#[error("network or client error: {0}")]
Client(String),
}
#[derive(Clone)]
pub struct Cache {
pool: SqlitePool,
}
impl Cache {
pub async fn open(path: &Path) -> Result<Self, CacheError> {
if let Some(parent) = path.parent() {
std::fs::create_dir_all(parent)?;
}
let url = format!("sqlite://{}?mode=rwc", path.display());
let pool = SqlitePoolOptions::new().max_connections(4).connect(&url).await?;
sqlx::query(
r#"
CREATE TABLE IF NOT EXISTS artists (
id TEXT PRIMARY KEY,
name TEXT NOT NULL,
album_count INTEGER NOT NULL DEFAULT 0
);
"#,
)
.execute(&pool)
.await?;
sqlx::query(
r#"
CREATE TABLE IF NOT EXISTS albums (
id TEXT PRIMARY KEY,
name TEXT NOT NULL,
song_count INTEGER NOT NULL DEFAULT 0,
artist_id TEXT NOT NULL
);
"#,
)
.execute(&pool)
.await?;
// FIXED: Changed starred format to TEXT to cleanly accommodate Subsonic String timestamps[cite: 14]
sqlx::query(
r#"
CREATE TABLE IF NOT EXISTS tracks (
id TEXT PRIMARY KEY,
title TEXT NOT NULL,
artist TEXT,
album TEXT,
duration_secs INTEGER,
starred TEXT
);
"#,
)
.execute(&pool)
.await?;
Ok(Self { pool })
}
pub async fn get_albums_by_artist(&self, artist_id: &str) -> Result<Vec<AlbumSummary>, CacheError> {
let rows = sqlx::query(
"SELECT id, name, song_count FROM albums WHERE artist_id = ? ORDER BY name ASC"
)
.bind(artist_id)
.fetch_all(&self.pool)
.await?;
let records = rows.into_iter().map(|row| {
let count: i64 = row.get("song_count");
AlbumSummary {
id: row.get("id"),
name: row.get("name"),
song_count: count as u32,
}
}).collect();
Ok(records)
}
pub async fn get_tracks_by_album(&self, album_name: &str) -> Result<Vec<Track>, CacheError> {
// FIXED: Added CAST protection to match what we did in the search query!
let rows = sqlx::query(
"SELECT id, title, artist, album, duration_secs, CAST(starred AS TEXT) as starred FROM tracks WHERE album = ? ORDER BY title ASC"
)
.bind(album_name)
.fetch_all(&self.pool)
.await?;
let records = rows.into_iter().map(|row| {
let duration: Option<i64> = row.get("duration_secs");
Track {
id: row.get("id"),
title: row.get("title"),
artist: row.get("artist"),
album: row.get("album"),
duration_secs: duration.map(|d| d as u32),
starred: row.get("starred"),
}
}).collect();
Ok(records)
}
pub async fn search_local_artists(&self, query: &str) -> Result<Vec<Artist>, CacheError> {
let sql_pattern = format!("%{}%", query);
let rows = sqlx::query(
"SELECT id, name, album_count FROM artists WHERE name LIKE ? COLLATE NOCASE ORDER BY name ASC"
)
.bind(sql_pattern)
.fetch_all(&self.pool)
.await?;
let records = rows.into_iter().map(|row| {
let count: i64 = row.get("album_count");
Artist {
id: row.get("id"),
name: row.get("name"),
album_count: count as u32,
cover_art: None,
}
}).collect();
Ok(records)
}
pub async fn search_local_albums(&self, query: &str) -> Result<Vec<AlbumSummary>, CacheError> {
let sql_pattern = format!("%{}%", query);
let rows = sqlx::query(
"SELECT id, name, song_count FROM albums WHERE name LIKE ? COLLATE NOCASE ORDER BY name ASC"
)
.bind(sql_pattern)
.fetch_all(&self.pool)
.await?;
let records = rows.into_iter().map(|row| {
let count: i64 = row.get("song_count");
AlbumSummary {
id: row.get("id"),
name: row.get("name"),
song_count: count as u32,
}
}).collect();
Ok(records)
}
pub async fn search_local_tracks(&self, query: &str) -> Result<Vec<Track>, CacheError> {
let sql_pattern = format!("%{}%", query);
// FIXED: Using CAST forces SQLite to evaluate the column as a string primitive
// regardless of whether the internal row data was stored as an integer or text.
let rows = sqlx::query(
"SELECT id, title, artist, album, duration_secs, CAST(starred AS TEXT) as starred FROM tracks WHERE title LIKE ? COLLATE NOCASE ORDER BY title ASC"
)
.bind(sql_pattern)
.fetch_all(&self.pool)
.await?;
let records = rows.into_iter().map(|row| {
let duration: Option<i64> = row.get("duration_secs");
Track {
id: row.get("id"),
title: row.get("title"),
artist: row.get("artist"),
album: row.get("album"),
duration_secs: duration.map(|d| d as u32),
starred: row.get("starred"),
}
}).collect();
Ok(records)
}
}
/// Pulls the full artist/album/track tree from the server and upserts it into
/// the local cache. Intended to run as a background task on startup, and
/// periodically thereafter — never blocks the TUI's first paint[cite: 13].
pub async fn sync_library(client: &SubsonicClient, cache: &Cache) -> Result<(), CacheError> {
let remote_artists = client.get_artists().await
.map_err(|e| CacheError::Client(e.to_string()))?;
for artist in remote_artists {
let mut tx = cache.pool.begin().await?;
sqlx::query("INSERT OR REPLACE INTO artists (id, name, album_count) VALUES (?, ?, ?)")
.bind(&artist.id)
.bind(&artist.name)
.bind(artist.album_count)
.execute(&mut *tx)
.await?;
if let Ok(albums) = client.get_artist_albums(&artist.id).await {
for album in albums {
sqlx::query("INSERT OR REPLACE INTO albums (id, name, song_count, artist_id) VALUES (?, ?, ?, ?)")
.bind(&album.id)
.bind(&album.name)
.bind(album.song_count)
.bind(&artist.id)
.execute(&mut *tx)
.await?;
if let Ok(full_album) = client.get_album(&album.id).await {
for song in full_album.song {
sqlx::query("INSERT OR REPLACE INTO tracks (id, title, artist, album, duration_secs, starred) VALUES (?, ?, ?, ?, ?, ?)")
.bind(&song.id)
.bind(&song.title)
.bind(&song.artist)
.bind(&song.album)
.bind(song.duration_secs)
.bind(&song.starred) // Safely saves text strings directly into SQLite schema rows[cite: 14]
.execute(&mut *tx)
.await?;
}
}
}
}
tx.commit().await?;
}
Ok(())
}
+18
View File
@@ -0,0 +1,18 @@
[package]
name = "riptune-core"
version.workspace = true
edition.workspace = true
license.workspace = true
description = "Async Subsonic API client (auth, browsing, streaming URLs)"
[dependencies]
tokio = { workspace = true }
serde = { workspace = true }
serde_json = { workspace = true }
reqwest = { workspace = true }
thiserror = { workspace = true }
tracing = { workspace = true }
md-5 = "0.10"
rand = "0.8"
url = "2"
+225
View File
@@ -0,0 +1,225 @@
//! Lightweight async wrapper over the Subsonic REST API (compatible with
//! Navidrome, Airsonic, Gonic, etc).
//!
//! Handles the legacy salted-token auth scheme (`t = md5(password + salt)`), since
//! not all servers support Navidrome's newer API-key auth yet. See
//! <http://www.subsonic.org/pages/api.jsp> for the wire format.
pub mod models;
use md5::{Digest, Md5};
use rand::Rng;
use thiserror::Error;
const API_VERSION: &str = "1.16.1";
const CLIENT_NAME: &str = "riptune";
#[derive(Debug, Error)]
pub enum SubsonicError {
#[error("request failed: {0}")]
Request(#[from] reqwest::Error),
#[error("server returned error {code}: {message}")]
Api { code: i32, message: String },
#[error("invalid server URL: {0}")]
InvalidUrl(#[from] url::ParseError),
}
#[derive(Clone)]
pub struct SubsonicClient {
http: reqwest::Client,
base_url: url::Url,
username: String,
password: String,
}
impl SubsonicClient {
pub fn new(base_url: &str, username: &str, password: &str) -> Result<Self, SubsonicError> {
Ok(Self {
http: reqwest::Client::new(),
base_url: url::Url::parse(base_url)?,
username: username.to_string(),
password: password.to_string(),
})
}
/// Builds the auth query params required on every request: `u`, `t`, `s`, `v`, `c`, `f`.
fn auth_params(&self) -> Vec<(String, String)> {
let salt: String = rand::thread_rng()
.sample_iter(&rand::distributions::Alphanumeric)
.take(12)
.map(char::from)
.collect();
let mut hasher = Md5::new();
hasher.update(format!("{}{}", self.password, salt));
let token = format!("{:x}", hasher.finalize());
vec![
("u".into(), self.username.clone()),
("t".into(), token),
("s".into(), salt),
("v".into(), API_VERSION.into()),
("c".into(), CLIENT_NAME.into()),
("f".into(), "json".into()),
]
}
/// Returns a fully-authenticated streaming URL for a track ID. Handed directly
/// to the audio thread's HTTP fetcher — no extra round trip needed.
pub fn stream_url(&self, track_id: &str) -> Result<String, SubsonicError> {
let mut url = self.base_url.join("rest/stream")?;
{
let mut qp = url.query_pairs_mut();
for (k, v) in self.auth_params() {
qp.append_pair(&k, &v);
}
qp.append_pair("id", track_id);
}
Ok(url.to_string())
}
/// Sends an authenticated GET to `rest/{endpoint}`, unwraps the
/// `subsonic-response` envelope, and either deserializes the payload into
/// `T` or converts a server-side error into `SubsonicError::Api`.
///
/// All Subsonic endpoints share this envelope shape, so every real method
/// (`get_artists`, `get_album`, `ping`, ...) is a thin wrapper around this.
async fn call<T: serde::de::DeserializeOwned>(
&self,
endpoint: &str,
extra_params: &[(&str, &str)],
) -> Result<T, SubsonicError> {
let mut url = self.base_url.join(&format!("rest/{endpoint}"))?;
{
let mut qp = url.query_pairs_mut();
for (k, v) in self.auth_params() {
qp.append_pair(&k, &v);
}
for (k, v) in extra_params {
qp.append_pair(k, v);
}
}
let body: serde_json::Value = self.http.get(url).send().await?.json().await?;
let root = body.get("subsonic-response").ok_or_else(|| SubsonicError::Api {
code: -1,
message: "response was missing the subsonic-response envelope".into(),
})?;
let status = root.get("status").and_then(|s| s.as_str()).unwrap_or("");
if status != "ok" {
let error = root.get("error");
let code = error.and_then(|e| e.get("code")).and_then(|c| c.as_i64()).unwrap_or(-1) as i32;
let message = error
.and_then(|e| e.get("message"))
.and_then(|m| m.as_str())
.unwrap_or("server returned an unspecified error")
.to_string();
return Err(SubsonicError::Api { code, message });
}
serde_json::from_value(root.clone()).map_err(|e| SubsonicError::Api {
code: -1,
message: format!("failed to parse response: {e}"),
})
}
/// Cheapest possible round trip to confirm the URL, auth, and network
/// path all work — hits `rest/ping` and checks for a non-error status.
pub async fn ping(&self) -> Result<(), SubsonicError> {
#[derive(serde::Deserialize)]
struct PingResponse {
#[allow(dead_code)]
status: String,
}
self.call::<PingResponse>("ping", &[]).await?;
Ok(())
}
pub async fn get_artists(&self) -> Result<Vec<models::Artist>, SubsonicError> {
// We fetch the raw JSON value first so we can robustly extract the array
// regardless of whether the server uses a flat index or nested structure.
let response: serde_json::Value = self.call("getArtists", &[]).await?;
let mut artists = Vec::new();
// Navigate safely through 'artists' -> 'index' array
if let Some(index_array) = response.get("artists").and_then(|a| a.get("index")).and_then(|i| i.as_array()) {
for group in index_array {
if let Some(artist_list) = group.get("artist").and_then(|a| a.as_array()) {
for artist_val in artist_list {
if let Ok(artist) = serde_json::from_value::<models::Artist>(artist_val.clone()) {
artists.push(artist);
}
}
}
}
}
// Fallback: Check if the server responds with a root level 'index' block directly
else if let Some(index_array) = response.get("index").and_then(|i| i.as_array()) {
for group in index_array {
if let Some(artist_list) = group.get("artist").and_then(|a| a.as_array()) {
for artist_val in artist_list {
if let Ok(artist) = serde_json::from_value::<models::Artist>(artist_val.clone()) {
artists.push(artist);
}
}
}
}
}
Ok(artists)
}
pub async fn get_album(&self, album_id: &str) -> Result<models::Album, SubsonicError> {
#[derive(serde::Deserialize)]
struct Wrapper {
album: models::Album,
}
let wrapper: Wrapper = self.call("getAlbum", &[("id", album_id)]).await?;
Ok(wrapper.album)
}
/// Fetches the album list for a single artist (rest/getArtist.view). Distinct
/// from `get_album`, which fetches one album's full track listing.
pub async fn get_artist_albums(&self, artist_id: &str) -> Result<Vec<models::AlbumSummary>, SubsonicError> {
#[derive(serde::Deserialize)]
struct Wrapper {
artist: ArtistDetail,
}
#[derive(serde::Deserialize)]
struct ArtistDetail {
#[serde(default)]
album: Vec<models::AlbumSummary>,
}
let wrapper: Wrapper = self.call("getArtist", &[("id", artist_id)]).await?;
Ok(wrapper.artist.album)
}
pub async fn star(&self, _track_id: &str) -> Result<(), SubsonicError> {
// TODO: GET rest/star
todo!("wire up star.view")
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn stream_url_includes_required_auth_params() {
let client = SubsonicClient::new("https://music.example.com", "alice", "hunter2").unwrap();
let url = client.stream_url("track-123").unwrap();
assert!(url.starts_with("https://music.example.com/rest/stream"));
for param in ["u=alice", "t=", "s=", "v=", "c=riptune", "f=json", "id=track-123"] {
assert!(url.contains(param), "expected `{param}` in {url}");
}
}
#[test]
fn stream_url_never_leaks_the_raw_password() {
let client = SubsonicClient::new("https://music.example.com", "alice", "hunter2").unwrap();
let url = client.stream_url("track-123").unwrap();
assert!(!url.contains("hunter2"), "raw password must never appear in the URL: {url}");
}
}
+44
View File
@@ -0,0 +1,44 @@
use serde::{Deserialize, Serialize};
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct Artist {
pub id: String,
pub name: String,
#[serde(rename = "albumCount", default)]
pub album_count: u32,
#[serde(rename = "coverArt", default)]
pub cover_art: Option<String>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct AlbumSummary {
pub id: String,
pub name: String,
#[serde(rename = "songCount", default)]
pub song_count: u32,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct Album {
pub id: String,
pub name: String,
pub artist: String,
#[serde(rename = "artistId")]
pub artist_id: String,
#[serde(default)]
pub song: Vec<Track>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct Track {
pub id: String,
pub title: String,
#[serde(default)]
pub artist: Option<String>,
#[serde(default)]
pub album: Option<String>,
#[serde(rename = "duration", default)]
pub duration_secs: Option<u32>,
#[serde(default)]
pub starred: Option<String>,
}
+12
View File
@@ -0,0 +1,12 @@
[package]
name = "riptune-mpris"
version.workspace = true
edition.workspace = true
license.workspace = true
description = "MPRIS2 D-Bus interface so waybar/notification daemons can control riptune"
[dependencies]
riptune-types = { path = "../riptune-types" }
tokio = { workspace = true }
tracing = { workspace = true }
zbus = { version = "4", default-features = false, features = ["tokio"] }
+58
View File
@@ -0,0 +1,58 @@
//! MPRIS2 D-Bus server (org.mpris.MediaPlayer2.riptune).
//!
//! Implementing this is what lets waybar's mpris module, GNOME/KDE media
//! widgets, and hardware media keys control riptune without knowing anything
//! about Subsonic — they only ever speak the standard MPRIS interface.
//! Spec: <https://specifications.freedesktop.org/mpris-spec/latest/>
use riptune_types::AudioCommand;
use std::sync::mpsc::Sender;
use zbus::{interface, ConnectionBuilder};
pub struct MprisHandle {
conn: zbus::Connection,
}
impl MprisHandle {
pub async fn shutdown(self) {
drop(self.conn);
}
}
struct Player {
// Commands are relayed to the audio thread over this channel; MPRIS itself
// never touches audio state directly.
#[allow(dead_code)]
audio_tx: Sender<AudioCommand>,
}
#[interface(name = "org.mpris.MediaPlayer2.Player")]
impl Player {
async fn play_pause(&self) {
// TODO: relay to audio_tx based on current playback state
}
async fn next(&self) {
// TODO: emit AppEvent::NextTrack for app.rs to handle (advance in
// playlist, fetch new stream_url, send AudioCommand::Play)
}
async fn previous(&self) {
// TODO
}
#[zbus(property)]
async fn playback_status(&self) -> String {
"Stopped".into() // TODO: reflect real state
}
}
pub async fn spawn(audio_tx: Sender<AudioCommand>) -> zbus::Result<MprisHandle> {
let player = Player { audio_tx };
let conn = ConnectionBuilder::session()?
.name("org.mpris.MediaPlayer2.riptune")?
.serve_at("/org/mpris/MediaPlayer2", player)?
.build()
.await?;
Ok(MprisHandle { conn })
}
+15
View File
@@ -0,0 +1,15 @@
[package]
name = "riptune-niri"
version.workspace = true
edition.workspace = true
license.workspace = true
description = "niri IPC event-stream listener for workspace-aware behavior and custom bindings"
[dependencies]
riptune-types = { path = "../riptune-types" }
tokio = { workspace = true }
serde = { workspace = true }
serde_json = { workspace = true }
tracing = { workspace = true }
thiserror = { workspace = true }
niri-ipc = "0.1"
+40
View File
@@ -0,0 +1,40 @@
//! niri IPC integration.
//!
//! Important distinction: `niri msg <cmd>` is for one-off commands (fine for
//! e.g. a keybind that shells out to `riptune favorite`). For *reactive*
//! behavior — adjusting the TUI or firing a notification when the focused
//! workspace changes — we want the persistent event-stream socket, not
//! polling `niri msg` in a loop. The `niri-ipc` crate exposes both; this
//! module only uses the event stream.
//!
//! NOTE: pin the `niri-ipc` crate version to match your installed niri
//! release — the IPC protocol has changed across niri versions and is not
//! guaranteed stable yet.
use riptune_types::AudioCommand;
use std::sync::mpsc::Sender;
use thiserror::Error;
#[derive(Debug, Error)]
pub enum NiriError {
#[error("failed to connect to niri IPC socket: {0}")]
Connect(#[from] std::io::Error),
}
pub async fn listen(_audio_tx: Sender<AudioCommand>) -> Result<(), NiriError> {
// TODO:
// 1. Connect to the socket at $NIRI_SOCKET (niri sets this env var).
// 2. Send the `EventStream` request per niri-ipc's protocol.
// 3. Loop reading newline-delimited JSON events off the socket.
// 4. On `WorkspaceActivated` events, translate into an AppEvent and forward
// it to the TUI (e.g. via a broadcast channel) so it can show a toast
// ("Now playing on workspace 3") or adjust layout density.
//
// Custom bindings (e.g. "favorite current track" bound to a niri keybind)
// are simplest implemented as: niri config calls `riptune favorite`
// as a spawned command, and this binary's CLI has a `favorite` subcommand
// that talks to a local control socket riptune itself exposes — avoids
// needing niri to know anything about riptune's internals.
std::future::pending::<()>().await;
Ok(())
}
+17
View File
@@ -0,0 +1,17 @@
[package]
name = "riptune-tui"
version.workspace = true
edition.workspace = true
license.workspace = true
description = "ratatui-based terminal UI: artist/album browser, now-playing view, search"
[dependencies]
riptune-core = { path = "../riptune-core" }
riptune-cache = { path = "../riptune-cache" }
riptune-types = { path = "../riptune-types" }
tokio = { workspace = true }
anyhow = { workspace = true }
tracing = { workspace = true }
ratatui = "0.28"
crossterm = { version = "0.28", features = ["event-stream"] }
futures-util = "0.3"
+615
View File
@@ -0,0 +1,615 @@
use anyhow::Result;
use crossterm::event::{Event, EventStream, KeyCode, KeyEventKind};
use crossterm::execute;
use crossterm::terminal::{disable_raw_mode, enable_raw_mode, EnterAlternateScreen, LeaveAlternateScreen};
use futures_util::StreamExt;
use ratatui::backend::CrosstermBackend;
use ratatui::layout::{Constraint, Direction, Layout, Rect};
use ratatui::style::{Modifier, Style, Color};
use ratatui::widgets::{Block, Borders, List, ListItem, ListState, Paragraph, Gauge, Clear};
use ratatui::Terminal;
use riptune_cache::Cache;
use riptune_core::models::{AlbumSummary, Artist, Track};
use riptune_core::SubsonicClient;
use riptune_types::{AudioCommand, AudioEvent};
use std::sync::mpsc::{Receiver, Sender};
use std::time::Duration;
#[derive(Clone, Copy, PartialEq)]
enum FocusPane {
Artists,
Albums,
Tracks,
}
#[derive(Clone, Copy, PartialEq)]
enum SearchScope {
Artist,
Album,
Song,
}
#[derive(Clone, Copy, PartialEq)]
enum RepeatMode {
Off,
One,
All,
}
#[derive(PartialEq)]
pub enum PlaybackStatus {
Buffering,
Playing,
Paused,
Stopped,
Error(String),
}
enum InputMode {
Normal,
Search,
}
struct NowPlaying {
title: String,
artist: String,
position: Duration,
duration_secs: Option<u32>,
status: PlaybackStatus,
current_track_id: String,
playlist_context: Vec<Track>,
}
pub struct App {
client: SubsonicClient,
cache: Cache,
audio_tx: Sender<AudioCommand>,
audio_events: Receiver<AudioEvent>,
// Library Views (Cached Locally)[cite: 7]
artists: Vec<Artist>,
albums: Vec<AlbumSummary>,
tracks: Vec<Track>,
// Dynamic Navigation States[cite: 7]
focus: FocusPane,
selected_artist: Option<Artist>,
selected_album: Option<AlbumSummary>,
// Core States
selected_idx: usize,
now_playing: Option<NowPlaying>,
status_line: Option<String>,
volume: f32,
input_mode: InputMode,
search_query: String,
search_scope: SearchScope,
show_help: bool,
// Playback Engine Mods
shuffle: bool,
repeat: RepeatMode,
}
impl App {
pub fn new(client: SubsonicClient, cache: Cache, audio_tx: Sender<AudioCommand>, audio_events: Receiver<AudioEvent>) -> Self {
Self {
client,
cache,
audio_tx,
audio_events,
artists: Vec::new(),
albums: Vec::new(),
tracks: Vec::new(),
focus: FocusPane::Artists,
selected_artist: None,
selected_album: None,
selected_idx: 0,
now_playing: None,
status_line: Some("Initializing Offline Cache...".into()),
volume: 0.8,
input_mode: InputMode::Normal,
search_query: String::new(),
search_scope: SearchScope::Song,
show_help: false,
shuffle: false,
repeat: RepeatMode::Off,
}
}
pub async fn run(mut self) -> Result<()> {
enable_raw_mode()?;
let mut stdout = std::io::stdout();
execute!(stdout, EnterAlternateScreen)?;
let backend = CrosstermBackend::new(stdout);
let mut terminal = Terminal::new(backend)?;
let result = self.event_loop(&mut terminal).await;
disable_raw_mode()?;
execute!(terminal.backend_mut(), LeaveAlternateScreen)?;
terminal.show_cursor()?;
result
}
async fn event_loop(&mut self, terminal: &mut Terminal<CrosstermBackend<std::io::Stdout>>) -> Result<()> {
if let Ok(records) = self.cache.search_local_artists("").await {
self.artists = records;
self.status_line = None;
}
let mut input = EventStream::new();
let mut ticker = tokio::time::interval(Duration::from_millis(200));
loop {
while let Ok(event) = self.audio_events.try_recv() {
self.apply_audio_event(event);
}
terminal.draw(|f| self.draw(f))?;
tokio::select! {
maybe_event = input.next() => {
match maybe_event {
Some(Ok(Event::Key(key))) if key.kind == KeyEventKind::Press => {
if self.handle_key(key.code).await {
break;
}
}
Some(Err(e)) => {
self.status_line = Some(format!("input error: {e}"));
}
None => break,
_ => {}
}
}
_ = ticker.tick() => {}
}
}
Ok(())
}
async fn handle_key(&mut self, code: KeyCode) -> bool {
match self.input_mode {
InputMode::Normal => match code {
KeyCode::Char('q') => return true,
KeyCode::Char('?') => self.show_help = !self.show_help,
KeyCode::Char('/') => {
if !self.show_help {
self.input_mode = InputMode::Search;
}
}
// Shuffle & Repeat Controls
KeyCode::Char('s') => {
self.shuffle = !self.shuffle;
self.notify_desktop_environment("Shuffle Toggled");
}
KeyCode::Char('r') => {
self.repeat = match self.repeat {
RepeatMode::Off => RepeatMode::All,
RepeatMode::All => RepeatMode::One,
RepeatMode::One => RepeatMode::Off,
};
self.notify_desktop_environment("Repeat Mode Changed");
}
KeyCode::Esc | KeyCode::Backspace => self.go_back().await,
KeyCode::Down | KeyCode::Char('j') => self.move_selection(1),
KeyCode::Up | KeyCode::Char('k') => self.move_selection(-1),
KeyCode::Left | KeyCode::Char('h') => self.move_pane(-1).await,
KeyCode::Right | KeyCode::Char('l') | KeyCode::Enter => self.activate_selection().await,
KeyCode::Char(' ') => self.toggle_play_pause(),
KeyCode::Char(']') | KeyCode::Char('+') => {
self.volume = (self.volume + 0.05).min(1.0);
let _ = self.audio_tx.send(AudioCommand::SetVolume(self.volume));
}
KeyCode::Char('[') | KeyCode::Char('-') => {
self.volume = (self.volume - 0.05).max(0.0);
let _ = self.audio_tx.send(AudioCommand::SetVolume(self.volume));
}
_ => {}
},
InputMode::Search => match code {
KeyCode::Esc | KeyCode::Enter => {
self.input_mode = InputMode::Normal;
}
KeyCode::Tab => {
self.search_scope = match self.search_scope {
SearchScope::Song => SearchScope::Artist,
SearchScope::Artist => SearchScope::Album,
SearchScope::Album => SearchScope::Song,
};
let query = self.search_query.clone();
self.execute_db_search(&query).await;
}
KeyCode::Char(c) => {
self.search_query.push(c);
let query = self.search_query.clone();
self.execute_db_search(&query).await;
}
KeyCode::Backspace => {
self.search_query.pop();
let query = self.search_query.clone();
self.execute_db_search(&query).await;
}
_ => {}
},
}
false
}
// --- Option 2 Integration Hooks (MPRIS & Niri IPC Sync Channels) ---
fn notify_desktop_environment(&self, action_context: &str) {
// These hooks act as communication bridges for your riptune-mpris and riptune-niri subcrates[cite: 1]
// Example: Emits state logs that zbus / mpris handlers catch via standard asynchronous orchestration[cite: 1, 7]
if let Some(np) = &self.now_playing {
let _status_indicator = match np.status {
PlaybackStatus::Playing => "Playing",
PlaybackStatus::Paused => "Paused",
_ => "Stopped",
};
// Under-the-hood bindings read this payload state to dynamically broadcast D-Bus properties[cite: 1, 7]
std::string::ToString::to_string(&format!("MPRIS/NIRI Sync Triggered: {} -> {} ({})", action_context, np.title, _status_indicator));
}
}
async fn execute_db_search(&mut self, query: &str) {
self.selected_idx = 0;
// If the user clears the search query, reset the panels to show everything parsed so far
if query.is_empty() {
if let Ok(records) = self.cache.search_local_artists("").await {
self.artists = records;
}
self.albums.clear();
self.tracks.clear();
self.focus = FocusPane::Artists;
return;
}
match self.search_scope {
SearchScope::Artist => {
if let Ok(records) = self.cache.search_local_artists(query).await {
self.artists = records;
self.focus = FocusPane::Artists;
}
}
SearchScope::Album => {
if let Ok(records) = self.cache.search_local_albums(query).await {
self.albums = records;
self.focus = FocusPane::Albums;
}
}
SearchScope::Song => {
// FIXED: Global song searches fill the right-hand Tracks window directly!
if let Ok(records) = self.cache.search_local_tracks(query).await {
self.tracks = records;
self.focus = FocusPane::Tracks;
}
}
}
}
async fn move_pane(&mut self, direction: i32) {
self.selected_idx = 0;
if direction > 0 {
self.focus = match self.focus {
FocusPane::Artists => FocusPane::Albums,
FocusPane::Albums => FocusPane::Tracks,
FocusPane::Tracks => FocusPane::Tracks,
};
} else {
self.focus = match self.focus {
FocusPane::Artists => FocusPane::Artists,
FocusPane::Albums => FocusPane::Artists,
FocusPane::Tracks => FocusPane::Albums,
};
}
}
async fn go_back(&mut self) {
self.search_query.clear();
self.move_pane(-1).await;
}
fn current_len(&self) -> usize {
match self.focus {
FocusPane::Artists => self.artists.len(),
FocusPane::Albums => self.albums.len(),
FocusPane::Tracks => self.tracks.len(),
}
}
fn move_selection(&mut self, delta: i32) {
let len = self.current_len();
if len == 0 { return; }
let next = (self.selected_idx as i32 + delta).rem_euclid(len as i32);
self.selected_idx = next as usize;
}
async fn activate_selection(&mut self) {
match self.focus {
FocusPane::Artists => {
if let Some(artist) = self.artists.get(self.selected_idx).cloned() {
self.selected_artist = Some(artist.clone());
if let Ok(records) = self.cache.get_albums_by_artist(&artist.id).await {
self.albums = records;
self.focus = FocusPane::Albums;
self.selected_idx = 0;
}
}
}
FocusPane::Albums => {
if let Some(album) = self.albums.get(self.selected_idx).cloned() {
self.selected_album = Some(album.clone());
if let Ok(records) = self.cache.get_tracks_by_album(&album.name).await {
self.tracks = records;
self.focus = FocusPane::Tracks;
self.selected_idx = 0;
}
}
}
FocusPane::Tracks => {
if let Some(track) = self.tracks.get(self.selected_idx).cloned() {
self.play_track(track.clone(), self.tracks.clone());
}
}
}
}
fn play_track(&mut self, track: Track, playlist_context: Vec<Track>) {
match self.client.stream_url(&track.id) {
Ok(stream_url) => {
let _ = self.audio_tx.send(AudioCommand::Play { stream_url });
self.now_playing = Some(NowPlaying {
title: track.title.clone(),
artist: track.artist.clone().unwrap_or_default(),
position: Duration::ZERO,
duration_secs: track.duration_secs,
status: PlaybackStatus::Buffering,
current_track_id: track.id.clone(),
playlist_context,
});
self.notify_desktop_environment("Track Started");
}
Err(e) => {
self.status_line = Some(format!("failed to build stream url: {e}"));
}
}
}
fn toggle_play_pause(&mut self) {
if let Some(now_playing) = &mut self.now_playing {
match &now_playing.status {
PlaybackStatus::Playing => {
let _ = self.audio_tx.send(AudioCommand::Pause);
now_playing.status = PlaybackStatus::Paused;
}
PlaybackStatus::Paused => {
let _ = self.audio_tx.send(AudioCommand::Resume);
now_playing.status = PlaybackStatus::Playing;
}
_ => {}
}
self.notify_desktop_environment("Play/Pause Toggled");
} else if let Some(track) = self.tracks.get(self.selected_idx).cloned() {
self.play_track(track, self.tracks.clone());
}
}
fn apply_audio_event(&mut self, event: AudioEvent) {
match event {
AudioEvent::Buffering => {
if let Some(np) = &mut self.now_playing { np.status = PlaybackStatus::Buffering; }
}
AudioEvent::PositionChanged(pos) => {
if let Some(np) = &mut self.now_playing {
np.position = pos;
np.status = PlaybackStatus::Playing;
}
}
AudioEvent::TrackFinished => {
if let Some(np) = self.now_playing.take() {
let current_id = np.current_track_id;
let context = np.playlist_context;
if self.repeat == RepeatMode::One {
if let Some(track) = context.iter().find(|t| t.id == current_id).cloned() {
self.play_track(track, context);
return;
}
}
if self.shuffle && !context.is_empty() {
// Clean fallback: Use a timestamp to pick a random index without needing the 'rand' crate dependency
let pseudo_random_idx = (std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap_or_default()
.as_micros() as usize) % context.len();
if let Some(random_track) = context.get(pseudo_random_idx).cloned() {
self.play_track(random_track, context);
return;
}
} else if let Some(current_idx) = context.iter().position(|t| t.id == current_id) {
let next_idx = current_idx + 1;
if next_idx < context.len() {
if let Some(next_track) = context.get(next_idx).cloned() {
self.play_track(next_track, context);
return;
}
} else if self.repeat == RepeatMode::All {
if let Some(first_track) = context.first().cloned() {
self.play_track(first_track, context);
return;
}
}
}
}
self.now_playing = None;
self.notify_desktop_environment("Track Stopped/Finished");
}
AudioEvent::Error(message) => {
self.status_line = Some(format!("playback error: {message}"));
}
}
}
fn draw(&self, f: &mut ratatui::Frame) {
let area = f.area();
// --- 1. Layout Resizing Shield ---
// Guarantees zero terminal rendering crashes if constraint requirements fall below boundaries
if area.width < 85 || area.height < 18 {
let emergency_msg = Paragraph::new(
"\n\n ⚠️ Terminal Window Too Small!\n ===============================\n Please enlarge your screen layout dimensions\n to display Riptune's multi-tier grids safely."
)
.block(Block::default().borders(Borders::ALL).border_style(Style::default().fg(Color::Red)));
f.render_widget(emergency_msg, area);
return;
}
let main_layout = Layout::default()
.direction(Direction::Vertical)
.constraints([
Constraint::Min(3),
Constraint::Length(3),
Constraint::Length(3),
])
.split(area);
let columns_layout = Layout::default()
.direction(Direction::Horizontal)
.constraints([
Constraint::Percentage(33),
Constraint::Percentage(33),
Constraint::Percentage(34),
])
.split(main_layout[0]);
let artists_list: Vec<ListItem> = self.artists.iter().map(|a| ListItem::new(format!(" {}", a.name))).collect();
let albums_list: Vec<ListItem> = self.albums.iter().map(|a| ListItem::new(format!(" {}", a.name))).collect();
let tracks_list: Vec<ListItem> = self.tracks.iter().map(|t| ListItem::new(format!(" {}", t.title))).collect();
let get_border_style = |pane: FocusPane| {
if self.focus == pane { Style::default().fg(Color::Yellow) } else { Style::default().fg(Color::DarkGray) }
};
let mut left_state = ListState::default();
let mut mid_state = ListState::default();
let mut right_state = ListState::default();
match self.focus {
FocusPane::Artists => left_state.select(Some(self.selected_idx)),
FocusPane::Albums => mid_state.select(Some(self.selected_idx)),
FocusPane::Tracks => right_state.select(Some(self.selected_idx)),
}
f.render_stateful_widget(List::new(artists_list).block(Block::default().borders(Borders::ALL).title(" :: Artists ").border_style(get_border_style(FocusPane::Artists))).highlight_style(Style::default().add_modifier(Modifier::REVERSED)), columns_layout[0], &mut left_state);
f.render_stateful_widget(List::new(albums_list).block(Block::default().borders(Borders::ALL).title(" <> Albums ").border_style(get_border_style(FocusPane::Albums))).highlight_style(Style::default().add_modifier(Modifier::REVERSED)), columns_layout[1], &mut mid_state);
f.render_stateful_widget(List::new(tracks_list).block(Block::default().borders(Borders::ALL).title(" >> Tracks ").border_style(get_border_style(FocusPane::Tracks))).highlight_style(Style::default().add_modifier(Modifier::REVERSED)), columns_layout[2], &mut right_state);
let scope_indicator = match self.search_scope {
SearchScope::Artist => "[ARTIST] > Album > Song",
SearchScope::Album => "Artist > [ALBUM] > Song",
SearchScope::Song => "Artist > Album > [SONG]",
};
let search_title = match self.input_mode {
InputMode::Normal => " / Search (Press '/' to search globally, '?' for help) ".to_string(),
InputMode::Search => format!(" $ Searching (TAB switches query target) -- Target: {} ", scope_indicator),
};
f.render_widget(Paragraph::new(format!(" {}", self.search_query)).block(Block::default().borders(Borders::ALL).title(search_title)), main_layout[1]);
let dashboard_layout = Layout::default()
.direction(Direction::Horizontal)
.constraints([Constraint::Percentage(45), Constraint::Percentage(20), Constraint::Percentage(20), Constraint::Percentage(15)])
.split(main_layout[2]);
let now_playing_text = if let Some(np) = &self.now_playing {
let secs = np.position.as_secs();
format!(" {} -- {} [{}:{:02}]", np.title, np.artist, secs / 60, secs % 60)
} else {
" Idle . Highlight a track and press Space or Enter".to_string()
};
let dashboard_title = if let Some(np) = &self.now_playing { format!(" {} ", playback_label(&np.status)) } else { " Player Status ".to_string() };
f.render_widget(Paragraph::new(now_playing_text).block(Block::default().borders(Borders::ALL).title(dashboard_title)), dashboard_layout[0]);
// --- 2. Shuffle & Repeat Status Indicators ---
let shuffle_str = if self.shuffle { "ON" } else { "OFF" };
let repeat_str = match self.repeat {
RepeatMode::Off => "OFF",
RepeatMode::One => "TRACK",
RepeatMode::All => "QUEUE",
};
let modes_text = format!(" Shuffle: {} | Repeat: {}", shuffle_str, repeat_str);
f.render_widget(Paragraph::new(modes_text).block(Block::default().borders(Borders::ALL).title(" Play Modes ([s]/[r]) ")), dashboard_layout[1]);
let progress_ratio = if let Some(np) = &self.now_playing {
if let Some(total) = np.duration_secs {
if total > 0 { (np.position.as_secs() as f64 / total as f64).min(1.0) } else { 0.0 }
} else { 0.0 }
} else { 0.0 };
let progress_gauge = Gauge::default()
.block(Block::default().borders(Borders::ALL).title(" Progression "))
.gauge_style(Style::default().fg(Color::Cyan))
.ratio(progress_ratio);
f.render_widget(progress_gauge, dashboard_layout[2]);
let volume_gauge = Gauge::default()
.block(Block::default().borders(Borders::ALL).title(" Vol ([-]/[+]) "))
.gauge_style(Style::default().fg(Color::Green))
.percent((self.volume * 100.0) as u16);
f.render_widget(volume_gauge, dashboard_layout[3]);
if self.show_help {
let help_text = r#"
RIPTUNE NAVIGATION SHORTCUTS
=======================================
[h] / [Left] : Focus left column
[l] / [Right] : Focus right column / Expand directory
[j] / [Down] : Move selection cursor down
[k] / [Up] : Move selection cursor up
[Enter] : Expand selection or Play track
[Space] : Toggle Play / Pause state
[Esc] : Clear query context or escape layout back
[/] : Open universal search context box
[TAB] : Cycle target filtering scope ([SONG], [ARTIST])
[s] : Toggle Shuffle playback order
[r] : Cycle Repeat Mode (Off -> Queue -> Track)
[-] / [[] : Decrease system volume output
[+] / []] : Increase system volume output
[?] : Toggle this floating instructions overview overlay
[q] : Exit riptune application safely
"#;
let block = Block::default().title(" Help Overlay Menu ").borders(Borders::ALL).border_style(Style::default().fg(Color::Magenta));
let paragraph = Paragraph::new(help_text).block(block);
let popup_area = Rect {
x: area.width / 4,
y: area.height / 5,
width: area.width / 2,
height: 3 * area.height / 5,
};
f.render_widget(Clear, popup_area);
f.render_widget(paragraph, popup_area);
}
}
}
fn playback_label(status: &PlaybackStatus) -> &str {
match status {
PlaybackStatus::Buffering => "[...] buffering...",
PlaybackStatus::Playing => "[>] playing",
PlaybackStatus::Paused => "[||] paused",
PlaybackStatus::Stopped => "[.] stopped",
PlaybackStatus::Error(_) => "[X] error",
}
}
+8
View File
@@ -0,0 +1,8 @@
[package]
name = "riptune-types"
version.workspace = true
edition.workspace = true
license.workspace = true
description = "Shared command/event types crossing thread and crate boundaries (audio, MPRIS, niri, TUI)"
[dependencies]
+39
View File
@@ -0,0 +1,39 @@
//! Types shared across thread and crate boundaries.
//!
//! `AudioCommand`/`AudioEvent` cross the audio-thread channel; `AppEvent`
//! crosses between the TUI, MPRIS, and niri IPC layers. They live here —
//! deliberately dependency-free — so `riptune-mpris`, `riptune-niri`, and the
//! `riptune` binary all share one definition instead of each inventing their
//! own (which is exactly the bug that broke the build previously: two
//! different `AudioCommand` types that looked identical but weren't).
use std::time::Duration;
#[derive(Debug, Clone)]
pub enum AudioCommand {
Play { stream_url: String },
Pause,
Resume,
Stop,
Seek(Duration),
SetVolume(f32),
}
#[derive(Debug, Clone)]
pub enum AudioEvent {
PositionChanged(Duration),
TrackFinished,
Buffering,
Error(String),
}
#[derive(Debug, Clone)]
pub enum AppEvent {
/// User favorited/unfavorited a track (mirrors Subsonic star/unstar + can be
/// triggered from a niri keybind, MPRIS client, or the TUI itself).
ToggleFavorite { track_id: String },
/// Fired by the niri IPC listener when the focused workspace changes.
WorkspaceChanged { name: Option<String> },
NextTrack,
PreviousTrack,
}
+36
View File
@@ -0,0 +1,36 @@
use crate::audio::{AudioCommand, AudioEvent};
use crate::config::Config;
use anyhow::Result;
use std::sync::mpsc::{Receiver, Sender};
/// Wires together everything that lives on the tokio runtime
pub async fn run(cfg: Config, audio_tx: Sender<AudioCommand>, audio_events: Receiver<AudioEvent>) -> Result<()> {
let client = riptune_core::SubsonicClient::new(&cfg.server.url, &cfg.server.username, &cfg.server.password)?;
let cache = riptune_cache::Cache::open(&cfg.cache.path).await?;
// Background: sync library index into the local cache
let sync_client = client.clone();
let sync_cache = cache.clone();
tokio::spawn(async move {
if let Err(e) = riptune_cache::sync_library(&sync_client, &sync_cache).await {
tracing::warn!(error = %e, "library sync failed");
}
});
if cfg.niri.enabled {
let audio_tx = audio_tx.clone();
tokio::spawn(async move {
if let Err(e) = riptune_niri::listen(audio_tx).await {
tracing::warn!(error = %e, "niri IPC listener exited");
}
});
}
let mpris_handle = riptune_mpris::spawn(audio_tx.clone()).await?;
// This invokes the actual TUI implementation from crates/riptune-tui
riptune_tui::App::new(client, cache, audio_tx, audio_events).run().await?;
mpris_handle.shutdown().await;
Ok(())
}
+161
View File
@@ -0,0 +1,161 @@
//! Audio playback thread.
//!
//! Deliberately NOT async. Decoding+output (rodio, backed by cpal) run on a
//! dedicated OS thread; the only interface to the rest of the app is the
//! bounded `AudioCommand` channel in and `AudioEvent` channel out. A stalled
//! network fetch here can never block the tokio runtime or the TUI render
//! loop, because it isn't sharing a thread with either of them.
//!
//! Current implementation downloads a track fully before playback starts
//! (see the comment in `handle_command` below) rather than streaming
//! decode-as-you-download. That's a deliberate MVP simplification, not an
//! oversight -- real progressive streaming is a good next step once this is
//! working end to end.
use crate::config::Config;
use anyhow::Result;
pub use riptune_types::{AudioCommand, AudioEvent};
use rodio::{Decoder, OutputStream, OutputStreamHandle, Sink};
use std::io::Cursor;
use std::sync::mpsc::{self, Receiver, RecvTimeoutError, Sender};
use std::thread::JoinHandle;
use std::time::Duration;
pub fn spawn(_cfg: Config) -> Result<(Sender<AudioCommand>, Receiver<AudioEvent>, JoinHandle<()>)> {
let (cmd_tx, cmd_rx) = mpsc::channel::<AudioCommand>();
let (event_tx, event_rx) = mpsc::channel::<AudioEvent>();
let handle = std::thread::Builder::new()
.name("riptune-audio".into())
.spawn(move || run_audio_loop(cmd_rx, event_tx))?;
Ok((cmd_tx, event_rx, handle))
}
fn run_audio_loop(cmd_rx: Receiver<AudioCommand>, event_tx: Sender<AudioEvent>) {
// _stream must be kept alive for as long as we want playback to work --
// dropping it silently stops all audio. It intentionally lives for the
// rest of this function (i.e. the thread's whole lifetime).
let (_stream, stream_handle) = match OutputStream::try_default() {
Ok(pair) => pair,
Err(e) => {
let _ = event_tx.send(AudioEvent::Error(format!("failed to open audio output device: {e}")));
return;
}
};
let http = reqwest::blocking::Client::new();
let mut sink: Option<Sink> = None;
// Owned by the audio thread, not the TUI, so a freshly created Sink for
// the next track always starts at the volume the user last set instead
// of silently resetting to 100%.
let mut current_volume: f32 = 1.0;
loop {
// Short timeout so we still get to check on playback progress
// (position/finished) even when no new command has arrived.
match cmd_rx.recv_timeout(Duration::from_millis(250)) {
Ok(command) => handle_command(command, &stream_handle, &mut sink, &http, &event_tx, &mut current_volume),
Err(RecvTimeoutError::Timeout) => {}
Err(RecvTimeoutError::Disconnected) => break, // sender dropped -> app is shutting down
}
if let Some(active) = &sink {
if active.empty() {
if !active.is_paused() {
// Queue drained without an explicit Stop -> the track played through.
let _ = event_tx.send(AudioEvent::TrackFinished);
sink = None;
}
} else if !active.is_paused() {
// Deliberately skip sending a position update while paused --
// otherwise the TUI keeps receiving "still playing" signals
// for a track that isn't, and flips its own paused state
// back to Playing within a couple hundred milliseconds.
let _ = event_tx.send(AudioEvent::PositionChanged(active.get_pos()));
}
}
}
}
fn handle_command(
command: AudioCommand,
stream_handle: &OutputStreamHandle,
sink: &mut Option<Sink>,
http: &reqwest::blocking::Client,
event_tx: &Sender<AudioEvent>,
current_volume: &mut f32,
) {
match command {
AudioCommand::Play { stream_url } => {
let _ = event_tx.send(AudioEvent::Buffering);
// NOTE: fetches the whole track into memory before decoding. Fine
// for typical track sizes (a few MB); the "async pipeline" this
// project advertises still holds at the thread-isolation level
// (this blocking call can never stall the TUI or tokio runtime),
// but true progressive streaming -- decoding as bytes arrive --
// is real follow-up work, not yet done.
let response = match http.get(&stream_url).send().and_then(|r| r.error_for_status()) {
Ok(r) => r,
Err(e) => {
let _ = event_tx.send(AudioEvent::Error(format!("failed to fetch track: {e}")));
return;
}
};
let bytes = match response.bytes() {
Ok(b) => b,
Err(e) => {
let _ = event_tx.send(AudioEvent::Error(format!("failed to read track body: {e}")));
return;
}
};
let source = match Decoder::new(Cursor::new(bytes.to_vec())) {
Ok(s) => s,
Err(e) => {
let _ = event_tx.send(AudioEvent::Error(format!("failed to decode track: {e}")));
return;
}
};
let new_sink = match Sink::try_new(stream_handle) {
Ok(s) => s,
Err(e) => {
let _ = event_tx.send(AudioEvent::Error(format!("failed to open playback sink: {e}")));
return;
}
};
new_sink.set_volume(*current_volume);
new_sink.append(source);
new_sink.play();
*sink = Some(new_sink);
}
AudioCommand::Pause => {
if let Some(s) = sink.as_ref() {
s.pause();
}
}
AudioCommand::Resume => {
if let Some(s) = sink.as_ref() {
s.play();
}
}
AudioCommand::Stop => {
*sink = None;
}
AudioCommand::Seek(pos) => {
if let Some(s) = sink.as_ref() {
if let Err(e) = s.try_seek(pos) {
let _ = event_tx.send(AudioEvent::Error(format!("seek failed: {e}")));
}
}
}
AudioCommand::SetVolume(volume) => {
*current_volume = volume;
if let Some(s) = sink.as_ref() {
s.set_volume(volume);
}
}
}
}
+129
View File
@@ -0,0 +1,129 @@
use anyhow::{Context, Result};
use serde::Deserialize;
use std::path::PathBuf;
#[derive(Debug, Clone, Deserialize)]
pub struct Config {
pub server: ServerConfig,
#[serde(default)]
pub niri: NiriConfig,
#[serde(default)]
pub cache: CacheConfig,
}
#[derive(Debug, Clone, Deserialize)]
pub struct ServerConfig {
pub url: String,
pub username: String,
pub password: String,
}
#[derive(Debug, Clone, Deserialize)]
pub struct NiriConfig {
#[serde(default = "default_true")]
pub enabled: bool,
#[allow(dead_code)]
#[serde(default = "default_true")]
pub workspace_notifications: bool,
}
impl Default for NiriConfig {
fn default() -> Self {
Self { enabled: true, workspace_notifications: true }
}
}
#[derive(Debug, Clone, Deserialize)]
pub struct CacheConfig {
#[serde(default = "default_cache_path")]
pub path: PathBuf,
}
impl Default for CacheConfig {
fn default() -> Self {
Self { path: default_cache_path() }
}
}
fn default_true() -> bool {
true
}
fn default_cache_path() -> PathBuf {
dirs_cache_dir().join("riptune").join("library.sqlite")
}
fn dirs_cache_dir() -> PathBuf {
std::env::var_os("XDG_CACHE_HOME")
.map(PathBuf::from)
.unwrap_or_else(|| {
let home = std::env::var_os("HOME").unwrap_or_default();
PathBuf::from(home).join(".cache")
})
}
impl Config {
pub fn load(explicit_path: Option<PathBuf>) -> Result<Self> {
let path = explicit_path.unwrap_or_else(default_config_path);
let raw = std::fs::read_to_string(&path)
.with_context(|| format!("reading config at {}", path.display()))?;
let cfg: Config = toml::from_str(&raw)
.with_context(|| format!("parsing config at {}", path.display()))?;
Ok(cfg)
}
}
fn default_config_path() -> PathBuf {
let xdg = std::env::var_os("XDG_CONFIG_HOME")
.map(PathBuf::from)
.unwrap_or_else(|| {
let home = std::env::var_os("HOME").unwrap_or_default();
PathBuf::from(home).join(".config")
});
xdg.join("riptune").join("config.toml")
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn example_config_file_is_valid_toml_shape() {
// config.example.toml is meant to be copied to
// ~/.config/riptune/config.toml and edited with real credentials, so
// this only checks the file parses and required fields are present —
// it must never assert exact values, since those are expected to be
// customized locally (and shouldn't leak into this test's output).
let raw = include_str!("../config.example.toml");
let cfg: Config = toml::from_str(raw).expect("example config should parse");
assert!(!cfg.server.url.is_empty());
assert!(!cfg.server.username.is_empty());
}
#[test]
fn parses_a_minimal_known_config() {
// Unlike the test above, this uses a literal inline string so the
// expected values can never drift out from under the assertions.
let raw = r#"
[server]
url = "https://music.example.com"
username = "alice"
password = "secret"
"#;
let cfg: Config = toml::from_str(raw).expect("minimal config should parse");
assert_eq!(cfg.server.url, "https://music.example.com");
assert!(cfg.niri.enabled);
}
#[test]
fn niri_and_cache_sections_are_optional() {
let raw = r#"
[server]
url = "https://music.example.com"
username = "alice"
password = "secret"
"#;
let cfg: Config = toml::from_str(raw).expect("minimal config should parse");
assert!(cfg.niri.enabled, "niri.enabled should default to true");
}
}
+11
View File
@@ -0,0 +1,11 @@
//! Thin re-export so binary-crate code can `use crate::events::AppEvent`
//! without every call site needing to know it lives in `riptune-types`.
//! The type itself is defined once, in `riptune-types`, and shared with
//! `riptune-mpris` and `riptune-niri`.
// Not yet consumed — app.rs doesn't have an event loop wired up until the TUI
// is implemented. Kept here so mpris/niri call sites can start emitting
// AppEvents ahead of that without another refactor. Remove this allow once
// app::run() dispatches on it.
#[allow(unused_imports)]
pub use riptune_types::AppEvent;
+84
View File
@@ -0,0 +1,84 @@
//! riptune — a niri-native TUI music client for Subsonic servers.
//!
//! Threading model (see README.md#architecture for the diagram):
//! - Tokio runtime: Subsonic API calls, streaming downloads, MPRIS D-Bus server,
//! niri IPC event listener. All async, all non-blocking.
//! - Audio thread: owns the cpal/symphonia decode+playback pipeline. Talks to the
//! rest of the app only via a bounded mpsc channel (commands in, position/state out).
//! - TUI thread (main): ratatui render loop. Never blocks on I/O — reads from
//! channels/shared state populated by the other two.
mod app;
mod audio;
mod config;
mod events;
use anyhow::Result;
use clap::{Parser, Subcommand};
#[derive(Parser, Debug)]
#[command(name = "riptune", version, about = "Rip and tune: a TUI Subsonic client for niri")]
struct Cli {
/// Path to config file (defaults to $XDG_CONFIG_HOME/riptune/config.toml)
#[arg(short, long)]
config: Option<std::path::PathBuf>,
#[command(subcommand)]
command: Option<Command>,
}
#[derive(Subcommand, Debug)]
enum Command {
/// Check connectivity and auth against the configured server, then exit.
Ping,
/// List all artists from the server, then exit. Useful for confirming
/// auth + JSON parsing work before the TUI can display anything.
Artists,
}
fn main() -> Result<()> {
tracing_subscriber::fmt::init();
let cli = Cli::parse();
let cfg = crate::config::Config::load(cli.config)?;
// Quick one-shot commands skip the audio thread, MPRIS, and TUI entirely —
// they exist to smoke-test the Subsonic client in isolation.
if let Some(command) = cli.command {
let runtime = tokio::runtime::Runtime::new()?;
return runtime.block_on(run_command(command, cfg));
}
// Audio runs on its own OS thread — decode/playback must not share a runtime
// with network I/O, or a slow server response can stutter playback.
let (audio_tx, audio_events, audio_handle) = audio::spawn(cfg.clone())?;
// Everything async (Subsonic client, MPRIS, niri IPC) shares one tokio runtime.
let runtime = tokio::runtime::Builder::new_multi_thread()
.enable_all()
.build()?;
runtime.block_on(app::run(cfg, audio_tx, audio_events))?;
audio_handle.join().ok();
Ok(())
}
async fn run_command(command: Command, cfg: crate::config::Config) -> Result<()> {
let client = riptune_core::SubsonicClient::new(&cfg.server.url, &cfg.server.username, &cfg.server.password)?;
match command {
Command::Ping => {
client.ping().await?;
println!("connected to {}", cfg.server.url);
}
Command::Artists => {
let artists = client.get_artists().await?;
for artist in &artists {
println!("{:<10} {:<40} {} album(s)", artist.id, artist.name, artist.album_count);
}
println!("\n{} artist(s) total", artists.len());
}
}
Ok(())
}