4 Commits
Author SHA1 Message Date
ab 7d6af59e97 fix lock 2026-09-07 12:41:27 +03:00
ab 32c976beb8 Added self updater 2026-09-07 12:36:48 +03:00
Ultradesu 2607314da3 Fix lock 2026-09-02 15:50:41 +01:00
Ultradesu f901836929 fix federation recovery after sleep 2026-09-02 15:10:20 +01:00
14 changed files with 1237 additions and 249 deletions
+6 -1
View File
@@ -98,7 +98,12 @@ jobs:
run: |
set -euo pipefail
gh release create "$RELEASE_TAG" artifacts/*/* \
# Hash archive bytes; manifest entries use the release asset basenames.
for archive in artifacts/*/*; do
(cd "$(dirname "$archive")" && sha256sum "$(basename "$archive")")
done > SHA256SUMS
gh release create "$RELEASE_TAG" artifacts/*/* SHA256SUMS \
--title "furumi ${RELEASE_TAG}" \
--notes "Release ${RELEASE_TAG}" \
--verify-tag
+7 -2
View File
@@ -7,6 +7,8 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
## [Unreleased]
## [0.2.8] - 2026-09-02
### Added
- Optional offline music-similarity search for local tracks, backed by
@@ -27,7 +29,7 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
### Changed
- Similarity wire types, bounds, validation, and stream framing now come from
the shared `music-dht 0.4.0` API so native, web, and future clients can
the shared `music-dht 0.4` API so native, web, and future clients can
interoperate without sharing an embedding implementation.
- Existing SQLite embeddings are backfilled once with compact 256-bit routing
signatures; new embeddings store them immediately without changing exact
@@ -41,6 +43,8 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
### Fixed
- Federation now recovers automatically after sleep, prolonged idle, or a
degraded rendezvous transport while preserving local playback and state.
- Current-track information (`Shift+I`) now uses the enriched queue entry, so
it shows the same complete metadata and similarity action as `I` on that
track in the queue.
@@ -74,5 +78,6 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
- 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
[Unreleased]: https://gt.hexor.cy/ab/furumi_tui/compare/v0.2.8...HEAD
[0.2.8]: https://gt.hexor.cy/ab/furumi_tui/compare/v0.2.7...v0.2.8
[0.2.5]: https://gt.hexor.cy/ab/furumi_tui/compare/v0.2.4...v0.2.5
Generated
+423 -235
View File
File diff suppressed because it is too large Load Diff
+13 -3
View File
@@ -1,6 +1,6 @@
[package]
name = "furumi_tui"
version = "0.2.7"
version = "0.2.9"
edition = "2024"
rust-version = "1.97"
description = "A federated P2P player for personal music libraries"
@@ -11,6 +11,16 @@ name = "furumi"
path = "src/main.rs"
[dependencies]
rustls = { version = "0.23", default-features = false, features = ["ring", "std", "tls12"] }
semver = "1"
self-replace = "1.5"
tempfile = "3"
fs2 = "0.4"
flate2 = "1"
tar = "0.4"
zip = { version = "4", default-features = false, features = ["deflate"] }
object = { version = "0.36", default-features = false, features = ["read", "std"] }
anyhow = "1.0.102"
blake3 = "1"
crokey = "1.4.0"
@@ -22,9 +32,9 @@ 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 = "0.4.0"
music-dht = "0.4.1"
ratatui = "0.30.1"
reqwest = { version = "0.12.28", default-features = false, features = ["rustls-tls", "stream"] }
reqwest = { version = "0.12.28", default-features = false, features = ["rustls-tls", "stream", "blocking"] }
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"] }
rusty-opus = "0.9.1"
+18
View File
@@ -124,6 +124,24 @@ Import a music directory from Furumi's command line:
Federation, trusted-device pairing, and key bindings are configured directly
inside the player.
### Manual updates
In **Settings → Updates**, select **Check for updates**, then **Install update**
when a newer stable GitHub release is available. Downloads run in the background.
After installation, restart `furumi` to use the new version; playback is not
restarted automatically. Wait for an active update operation to finish before
quitting.
Updates replace the running executable in its installation directory, which
must be writable by your user. Release archives must include a matching entry
in the release's `SHA256SUMS` asset. Older releases without it cannot be installed
through this feature. The updater checks SHA-256 and the executable's format
and architecture before replacing it. Checksums provide integrity checking,
not publisher signatures. Settings and the local library are preserved.
There are no automatic startup checks. Only the existing release asset naming
scheme is supported; missing or incompatible platform builds are rejected.
### Now playing in tmux
While Furumi is running, a second invocation can print a cheap, single-line
+3
View File
@@ -10,6 +10,9 @@ use crate::library::models::{
/// the playback engine, imports). Tasks never touch AppState directly.
#[derive(Debug)]
pub enum AppEvent {
UpdateChecked(Result<Option<crate::updater::Update>, String>),
UpdateProgress(String),
UpdateInstalled(Result<(), String>),
StatusMessage(String),
/// A page of the artists list arrived (or failed).
ArtistsLoaded(Result<ArtistsPage, String>),
+64
View File
@@ -312,6 +312,7 @@ pub async fn run(
Arc::clone(&similarity),
settings.music_dir.clone(),
);
federation.start_supervisor();
state.music_dir = federation.media_dir();
state.federation.settings = federation.settings();
state.federation.devices = Some(devices.status());
@@ -430,6 +431,13 @@ pub async fn run(
}
if state.should_quit {
// A replacement must finish before the runtime/process is torn down.
if state.updater.busy {
state.should_quit = false;
state.status_message =
Some("Please wait for the update operation to finish".into());
continue;
}
if runtime.plain_text_mode {
leave_plain_text_mode()?;
runtime.plain_text_mode = false;
@@ -1414,6 +1422,31 @@ fn perform_effect(state: &mut AppState, runtime: &mut Runtime, effect: Effect) {
return;
}
match effect {
Effect::CheckUpdate => {
let tx = runtime.event_tx.clone();
tokio::spawn(async move {
let result = tokio::task::spawn_blocking(crate::updater::check)
.await
.map_err(|err| format!("update worker failed: {err}"))
.and_then(|result| result.map_err(|err| format!("{err:#}")));
let _ = tx.send(AppEvent::UpdateChecked(result));
});
}
Effect::InstallUpdate(update) => {
let tx = runtime.event_tx.clone();
tokio::spawn(async move {
let progress_tx = tx.clone();
let result = tokio::task::spawn_blocking(move || {
crate::updater::install(&update, |message| {
let _ = progress_tx.send(AppEvent::UpdateProgress(message));
})
})
.await
.map_err(|err| format!("update worker failed: {err}"))
.and_then(|result| result.map_err(|err| format!("{err:#}")));
let _ = tx.send(AppEvent::UpdateInstalled(result));
});
}
Effect::PlayCurrent => {
play_current(state, runtime);
push_media_metadata(state, runtime);
@@ -3097,6 +3130,37 @@ fn handle_playback_command(
fn handle_app_event(state: &mut AppState, runtime: &mut Runtime, event: AppEvent) {
match event {
AppEvent::UpdateChecked(result) => {
state.updater.busy = false;
state.updater.message = match result {
Ok(Some(update)) => {
let message = format!("v{} available", update.version);
state.updater.available = Some(update);
message
}
Ok(None) => "No newer stable release".into(),
Err(error) => format!("Check failed: {error}"),
};
state.status_message = Some(state.updater.message.clone());
}
AppEvent::UpdateProgress(message) => state.updater.message = message,
AppEvent::UpdateInstalled(result) => {
state.updater.busy = false;
state.updater.message = match result {
Ok(()) => {
state.updater.installed = true;
let version = state
.updater
.available
.as_ref()
.map(|u| u.version.as_str())
.unwrap_or("new version");
format!("Installed {version}; restart furumi")
}
Err(error) => format!("Update failed: {error}"),
};
state.status_message = Some(state.updater.message.clone());
}
AppEvent::StatusMessage(message) => state.status_message = Some(message),
AppEvent::ListenHistoryLoaded(result) => {
state.listen_history = Some(match result {
+8 -1
View File
@@ -926,6 +926,8 @@ impl FedRow {
/// the config directory.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum SettingsRow {
CheckUpdate,
InstallUpdate,
MusicDirectory,
Similarity(SimilarityRow),
Federation(FedRow),
@@ -1065,7 +1067,11 @@ pub fn device_status_order(state: &AppState) -> Vec<usize> {
}
pub fn settings_rows(state: &AppState) -> Vec<SettingsRow> {
let mut rows = vec![SettingsRow::MusicDirectory];
let mut rows = vec![
SettingsRow::MusicDirectory,
SettingsRow::CheckUpdate,
SettingsRow::InstallUpdate,
];
rows.extend(SimilarityRow::ALL.into_iter().map(SettingsRow::Similarity));
rows.extend(FedRow::ALL.into_iter().map(SettingsRow::Federation));
rows.push(SettingsRow::DeviceName);
@@ -1571,6 +1577,7 @@ impl DevicePlaybackState {
/// event handlers in the main loop; views render from `&AppState`.
#[derive(Debug, Default)]
pub struct AppState {
pub updater: crate::updater::State,
pub active_tab: Tab,
pub should_quit: bool,
pub shutting_down: bool,
+20
View File
@@ -18,6 +18,8 @@ pub const QUIT_CONFIRM_HINT: &str = "press quit again to exit";
/// owns the Runtime (audio controller, API client). Keeps update() pure.
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum Effect {
CheckUpdate,
InstallUpdate(crate::updater::Update),
/// (Re)start playback of `queue[queue_pos]`.
PlayCurrent,
TogglePause,
@@ -2815,6 +2817,24 @@ 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, SimilarityRow};
match settings_rows(state).get(state.settings_cursor).copied()? {
SettingsRow::CheckUpdate => {
if state.updater.busy || state.updater.installed {
return None;
}
state.updater.busy = true;
state.updater.available = None;
state.updater.message = "Checking GitHub Releases...".into();
Some(Effect::CheckUpdate)
}
SettingsRow::InstallUpdate => {
if state.updater.busy || state.updater.installed {
return None;
}
let update = state.updater.available.clone()?;
state.updater.busy = true;
state.updater.message = "Downloading update...".into();
Some(Effect::InstallUpdate(update))
}
SettingsRow::MusicDirectory => {
if state.music_dir_changing {
state.status_message = Some("music directory change is already running".into());
+21
View File
@@ -1,6 +1,27 @@
use super::*;
use crate::library::models::{ArtistCard, ArtistDetail, TrackItem};
#[test]
fn manual_update_check_is_single_flight_and_disabled_after_install() {
let mut state = AppState::default();
state.settings_cursor = settings_rows(&state)
.iter()
.position(|row| *row == crate::app::state::SettingsRow::CheckUpdate)
.unwrap();
assert_eq!(federation_select(&mut state), Some(Effect::CheckUpdate));
assert!(state.updater.busy);
assert_eq!(federation_select(&mut state), None);
state.updater.busy = false;
state.updater.installed = true;
assert_eq!(federation_select(&mut state), None);
state.updater.installed = false;
state.settings_cursor = settings_rows(&state)
.iter()
.position(|row| *row == crate::app::state::SettingsRow::InstallUpdate)
.unwrap();
assert_eq!(federation_select(&mut state), None);
}
fn with_artists(n: usize) -> AppState {
let mut state = AppState::default();
state.global.artists = (0..n)
+121 -4
View File
@@ -21,8 +21,8 @@ use std::collections::{HashMap, VecDeque};
use std::path::{Path, PathBuf};
use std::str::FromStr;
use std::sync::Arc;
use std::sync::atomic::{AtomicI64, Ordering};
use std::time::Duration;
use std::sync::atomic::{AtomicBool, AtomicI64, AtomicU64, Ordering};
use std::time::{Duration, Instant};
use anyhow::{Context, Result};
use music_dht::similarity_dht::SimilarityDht;
@@ -46,6 +46,11 @@ pub use similarity::SIMILARITY_ALPN;
/// How often the published library is re-synchronized with the local index.
const SYNC_INTERVAL: Duration = Duration::from_secs(60);
/// How often the application checks whether the shared transport needs a
/// full restart so all application-owned protocol acceptors are recreated.
const SUPERVISOR_INTERVAL: Duration = Duration::from_secs(15);
/// Prevents repeated restarts during a prolonged external network outage.
const RECOVERY_COOLDOWN: Duration = Duration::from_secs(5 * 60);
/// How many times a share-link content lookup is retried before the label
/// fallback kicks in.
@@ -370,6 +375,12 @@ pub struct FedStatus {
pub endpoint_id: String,
pub dht_node_id: String,
pub connected_peers: Vec<String>,
pub network_health: String,
pub rendezvous_failures: u32,
pub peer_dial_failures: u32,
pub rendezvous_restarts: u64,
pub recovery_count: u64,
pub last_rendezvous_error: Option<String>,
pub known_contacts: usize,
pub stored_dht_records: Option<usize>,
pub stored_dht_bytes: Option<u64>,
@@ -419,6 +430,10 @@ pub struct Federation {
metadata_cache: std::sync::Mutex<std::collections::HashMap<String, CachedTrackMetadata>>,
settings: std::sync::Mutex<FedSettings>,
running: tokio::sync::Mutex<Option<Running>>,
supervisor_task: std::sync::Mutex<Option<tokio::task::JoinHandle<()>>>,
supervisor_shutdown: tokio::sync::Notify,
shutting_down: AtomicBool,
recovery_count: AtomicU64,
last_sync: std::sync::Mutex<Option<String>>,
last_error: std::sync::Mutex<Option<String>>,
transport_stats: Arc<TransportStats>,
@@ -549,6 +564,10 @@ impl Federation {
metadata_cache: std::sync::Mutex::new(Default::default()),
settings: std::sync::Mutex::new(load_settings()),
running: tokio::sync::Mutex::new(None),
supervisor_task: std::sync::Mutex::new(None),
supervisor_shutdown: tokio::sync::Notify::new(),
shutting_down: AtomicBool::new(false),
recovery_count: AtomicU64::new(0),
last_sync: std::sync::Mutex::new(None),
last_error: std::sync::Mutex::new(initial_error),
transport_stats: Arc::new(TransportStats::default()),
@@ -560,6 +579,18 @@ impl Federation {
lock(&self.settings).clone()
}
/// Starts the application-level federation supervisor once.
pub fn start_supervisor(self: &Arc<Self>) {
let mut task = lock(&self.supervisor_task);
if task.is_some() || self.shutting_down.load(Ordering::SeqCst) {
return;
}
let federation = Arc::clone(self);
*task = Some(tokio::spawn(async move {
federation.supervisor_loop().await;
}));
}
pub fn media_dir(&self) -> PathBuf {
lock(&self.media_dir).clone()
}
@@ -665,12 +696,43 @@ impl Federation {
network_id: NetworkId,
network_name: String,
) -> Result<()> {
self.start_with_network_id_mode(network_id, network_name, false)
.await
.map(|_| ())
}
async fn start_with_network_id_mode(
self: &Arc<Self>,
network_id: NetworkId,
network_name: String,
recovery_only: bool,
) -> Result<bool> {
let mut guard = self.running.lock().await;
if recovery_only {
let current = self.settings();
if !current.enabled
|| current.network_id.trim() != network_name
|| NetworkId::from_name(current.network_id.trim()) != network_id
{
return Ok(false);
}
}
if let Some(running) = guard.as_ref() {
if running.network_id == network_id {
return Ok(());
if !recovery_only || !running.service.network_health().restart_recommended {
return Ok(false);
}
tracing::warn!(
network = %network_name,
health = %running.service.network_health().state,
"restarting degraded federation service"
);
} else if recovery_only {
return Ok(false);
}
stop_running(guard.take()).await;
} else if recovery_only {
tracing::warn!(network = %network_name, "retrying stopped federation service");
}
std::fs::create_dir_all(&self.data_dir)
.with_context(|| format!("creating {}", self.data_dir.display()))?;
@@ -828,7 +890,7 @@ impl Federation {
],
});
self.set_error(None);
Ok(())
Ok(true)
}
async fn stop(&self) {
@@ -837,9 +899,57 @@ impl Federation {
}
pub async fn shutdown(&self) {
self.shutting_down.store(true, Ordering::SeqCst);
self.supervisor_shutdown.notify_one();
let supervisor = lock(&self.supervisor_task).take();
if let Some(supervisor) = supervisor {
let _ = supervisor.await;
}
self.stop().await;
}
async fn supervisor_loop(self: Arc<Self>) {
let mut interval = tokio::time::interval(SUPERVISOR_INTERVAL);
interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
interval.tick().await;
let mut last_attempt = None;
loop {
tokio::select! {
_ = self.supervisor_shutdown.notified() => return,
_ = interval.tick() => {}
}
if self.shutting_down.load(Ordering::SeqCst) {
return;
}
if last_attempt.is_some_and(|attempt: Instant| attempt.elapsed() < RECOVERY_COOLDOWN) {
continue;
}
match self.recover_if_needed().await {
Ok(false) => {}
Ok(true) => {
last_attempt = Some(Instant::now());
self.recovery_count.fetch_add(1, Ordering::Relaxed);
self.spawn_sync_soon().await;
}
Err(err) => {
last_attempt = Some(Instant::now());
tracing::error!(error = %err, "federation recovery failed");
self.set_error(Some(format!("network recovery failed: {err}")));
}
}
}
}
async fn recover_if_needed(self: &Arc<Self>) -> Result<bool> {
let settings = self.settings();
if !settings.enabled || settings.network_id.trim().is_empty() {
return Ok(false);
}
let network_name = settings.network_id.trim().to_string();
self.start_with_network_id_mode(NetworkId::from_name(&network_name), network_name, true)
.await
}
async fn service(&self) -> Result<Arc<MusicDhtService>> {
self.running
.lock()
@@ -975,6 +1085,7 @@ impl Federation {
let guard = self.running.lock().await;
let mut status = FedStatus {
network: settings.network_id,
recovery_count: self.recovery_count.load(Ordering::Relaxed),
last_sync: lock(&self.last_sync).clone(),
last_error: lock(&self.last_error).clone(),
protocols: ProtocolVersions::snapshot(&self.observed_protocols),
@@ -982,6 +1093,7 @@ impl Federation {
};
if let Some(running) = guard.as_ref() {
let service = &running.service;
let health = service.network_health();
status.running = true;
status.network = running.network_name.clone();
status.endpoint_id = service.endpoint_id().to_string();
@@ -991,6 +1103,11 @@ impl Federation {
.iter()
.map(|p| p.to_string())
.collect();
status.network_health = health.state.to_string();
status.rendezvous_failures = health.consecutive_rendezvous_failures;
status.peer_dial_failures = health.consecutive_peer_dial_failures;
status.rendezvous_restarts = health.rendezvous_restarts;
status.last_rendezvous_error = health.last_rendezvous_error;
status.known_contacts = service.known_peers().len();
status.stored_dht_records = service.dht_record_count().await.ok();
status.stored_dht_bytes =
+1
View File
@@ -12,6 +12,7 @@ mod similarity;
mod status;
mod streaming;
mod ui;
mod updater;
mod visualizer;
use std::io;
+100 -3
View File
@@ -4,7 +4,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};
use ratatui::widgets::{Block, Paragraph, Wrap};
use super::theme;
use crate::app::state::{AppState, DevicePresenceSection, FedRow, SimilarityRow, settings_rows};
@@ -33,7 +33,7 @@ pub fn draw(frame: &mut Frame, area: Rect, state: &AppState) {
}
let rows_height =
(settings_rows(state).len() + 10 + device_presence_sections(state).len()) as u16;
(settings_rows(state).len() + 15 + device_presence_sections(state).len()) as u16;
let [rows_area, _, status_area] = Layout::vertical([
Constraint::Length(rows_height.min(inner.height)),
Constraint::Length(1),
@@ -88,6 +88,48 @@ fn draw_settings_rows(frame: &mut Frame, area: Rect, state: &AppState) {
y = y.saturating_add(1);
draw_section(frame, area, state, &mut y, "Updates");
draw_row_enabled(
frame,
area,
state,
&mut y,
cursor,
state.settings_cursor,
"Check for updates",
format!("v{} | enter", env!("CARGO_PKG_VERSION")),
!state.updater.busy && !state.updater.installed,
);
cursor += 1;
draw_row_enabled(
frame,
area,
state,
&mut y,
cursor,
state.settings_cursor,
"Install update",
state
.updater
.available
.as_ref()
.map(|update| format!("v{} | enter", update.version))
.unwrap_or_else(|| "check for updates first".into()),
state.updater.available.is_some() && !state.updater.busy && !state.updater.installed,
);
cursor += 1;
if y < area.bottom() {
let height = 3.min(area.bottom() - y);
frame.render_widget(
Paragraph::new(state.updater.message.as_str())
.style(theme::dim())
.wrap(Wrap { trim: true }),
Rect::new(area.x, y, area.width, height),
);
y += height;
}
y = y.saturating_add(1);
draw_section(frame, area, state, &mut y, "Similarity Search");
let similarity = &state.similarity.settings;
for row in SimilarityRow::ALL {
@@ -961,7 +1003,10 @@ fn node_summary_lines(state: &AppState) -> Vec<Line<'static>> {
]
}
Some(status) => vec![
summary_line("Node", format!("running on {}", status.network)),
summary_line(
"Node",
format!("running on {} · {}", status.network, status.network_health),
),
summary_line(
"Peers",
format!(
@@ -1215,6 +1260,17 @@ fn status_detail_status_lines(state: &AppState, status_cursor: usize) -> Vec<Lin
status.known_contacts
),
));
lines.push(status_line(
"Discovery",
format!(
"{} · {} DHT / {} dial failures · {} rebuilds · {} full recoveries",
status.network_health,
status.rendezvous_failures,
status.peer_dial_failures,
status.rendezvous_restarts,
status.recovery_count
),
));
if !status.connected_peers.is_empty() {
let mut peers: Vec<String> = status
.connected_peers
@@ -1260,6 +1316,9 @@ fn status_detail_status_lines(state: &AppState, status_cursor: usize) -> Vec<Lin
if let Some(error) = &status.last_error {
lines.push(status_line("Error", first_line(error)));
}
if let Some(error) = &status.last_rendezvous_error {
lines.push(status_line("Discovery error", first_line(error)));
}
push_transport_summary_status(&mut lines, state, status);
}
}
@@ -1561,3 +1620,41 @@ fn relative_time_label(value_ms: Option<i64>, now_ms: i64) -> String {
format!("{}d ago", seconds / 60 / 60 / 24)
}
}
#[cfg(test)]
mod update_ui_tests {
use super::*;
#[test]
fn update_controls_and_existing_sections_render_in_both_layouts() {
for width in [80, 160] {
let mut terminal =
ratatui::Terminal::new(ratatui::backend::TestBackend::new(width, 55)).unwrap();
let mut state = AppState::default();
state.updater.message = "No newer stable release".into();
terminal
.draw(|frame| draw(frame, frame.area(), &state))
.unwrap();
let text: String = terminal
.backend()
.buffer()
.content
.iter()
.map(|cell| cell.symbol())
.collect();
for expected in [
"Music save directory",
"Check for updates",
"Install update",
"No newer stable release",
"Similarity Search",
"Federation",
] {
assert!(
text.contains(expected),
"missing {expected} at width {width}"
);
}
}
}
}
+432
View File
@@ -0,0 +1,432 @@
//! Manual portable-binary updates. All network/disk work runs off the UI thread.
use std::{
fs::File,
io::{Read, Write},
path::Path,
time::Duration,
};
use anyhow::{Context, Result, bail, ensure};
use object::Object;
use reqwest::blocking::Client;
use semver::Version;
use serde::Deserialize;
use sha2::{Digest, Sha256};
const REPOSITORY: &str = "https://api.github.com/repos/house-of-vanity/furumi_tui";
const MAX_ARCHIVE: u64 = 512 * 1024 * 1024;
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct Update {
pub version: String,
asset: Asset,
checksums: Asset,
}
#[derive(Debug, Clone, PartialEq, Eq, Deserialize)]
struct Asset {
name: String,
browser_download_url: String,
}
#[derive(Deserialize)]
struct Release {
tag_name: String,
draft: bool,
prerelease: bool,
assets: Vec<Asset>,
}
#[derive(Debug, Default)]
pub struct State {
pub busy: bool,
pub installed: bool,
pub available: Option<Update>,
pub message: String,
}
fn client() -> Result<Client> {
// Reuse the ring backend already used by federation; respect an existing provider.
let _ = rustls::crypto::ring::default_provider().install_default();
Ok(Client::builder()
.user_agent(concat!("furumi/", env!("CARGO_PKG_VERSION")))
.https_only(true)
.connect_timeout(Duration::from_secs(10))
.timeout(Duration::from_secs(300))
.build()?)
}
fn read_limited(mut reader: impl Read, limit: u64) -> Result<Vec<u8>> {
let mut data = Vec::new();
reader.by_ref().take(limit + 1).read_to_end(&mut data)?;
ensure!(data.len() as u64 <= limit, "download exceeds size limit");
Ok(data)
}
fn select_release(release: Release, current: &str, os: &str, arch: &str) -> Result<Option<Update>> {
let version = release
.tag_name
.strip_prefix('v')
.unwrap_or(&release.tag_name);
let next = Version::parse(version).context("invalid release version")?;
if release.draft
|| release.prerelease
|| !next.pre.is_empty()
|| next <= Version::parse(current)?
{
return Ok(None);
}
let platform = match os {
"linux" => "linux",
"macos" => "macos",
"windows" => "windows",
_ => bail!("self-update is unsupported on {os}"),
};
let extension = if os == "windows" { "zip" } else { "tar.gz" };
let name = format!("furumi-{platform}-{arch}-{version}.{extension}");
let find = |name: &str| -> Result<Asset> {
let matches: Vec<_> = release
.assets
.iter()
.filter(|asset| asset.name == name)
.collect();
ensure!(matches.len() == 1, "release has no unique {name} asset");
Ok(matches[0].clone())
};
Ok(Some(Update {
version: version.to_owned(),
asset: find(&name)?,
checksums: find("SHA256SUMS")?,
}))
}
pub fn check() -> Result<Option<Update>> {
let response = client()?
.get(format!("{REPOSITORY}/releases/latest"))
.timeout(Duration::from_secs(20))
.send()?;
if response.status() == reqwest::StatusCode::NOT_FOUND {
return Ok(None);
}
let release = serde_json::from_slice(&read_limited(
response.error_for_status()?,
2 * 1024 * 1024,
)?)?;
select_release(
release,
env!("CARGO_PKG_VERSION"),
std::env::consts::OS,
std::env::consts::ARCH,
)
}
fn checksum(text: &str, name: &str) -> Result<String> {
let mut found = None;
for line in text.lines() {
let Some((hash, filename)) = line.split_once(' ') else {
continue;
};
if filename.trim_start().trim_start_matches('*') != name {
continue;
}
ensure!(found.is_none(), "duplicate checksum for {name}");
ensure!(
hash.len() == 64 && hash.bytes().all(|b| b.is_ascii_hexdigit()),
"invalid SHA-256 for {name}"
);
found = Some(hash.to_ascii_lowercase());
}
found.context("release does not contain a checksum for the selected archive")
}
fn extract(archive: &Path, name: &str, output: &mut File) -> Result<()> {
let binary = if cfg!(windows) {
"furumi.exe"
} else {
"furumi"
};
let root = name
.strip_suffix(".tar.gz")
.or_else(|| name.strip_suffix(".zip"))
.context("unsupported archive")?;
// Existing release archives omit the version in their inner directory.
let root = root.rsplit_once('-').context("invalid archive name")?.0;
let expected = format!("{root}/{binary}");
let mut count = 0;
if name.ends_with(".zip") {
let mut archive = zip::ZipArchive::new(File::open(archive)?)?;
for index in 0..archive.len() {
let mut entry = archive.by_index(index)?;
if entry.name() != expected {
continue;
}
ensure!(
entry.is_file() && !entry.is_symlink(),
"binary is not a regular file"
);
count += 1;
ensure!(count == 1, "duplicate binary in archive");
output.write_all(&read_limited(&mut entry, MAX_ARCHIVE)?)?;
}
} else {
let mut archive = tar::Archive::new(flate2::read::GzDecoder::new(File::open(archive)?));
for entry in archive.entries()? {
let mut entry = entry?;
if entry.path()?.as_ref() != Path::new(&expected) {
continue;
}
ensure!(
entry.header().entry_type().is_file(),
"binary is not a regular file"
);
count += 1;
ensure!(count == 1, "duplicate binary in archive");
output.write_all(&read_limited(&mut entry, MAX_ARCHIVE)?)?;
}
}
ensure!(count == 1, "archive does not contain {expected}");
output.sync_all()?;
Ok(())
}
fn validate_binary(current: &[u8], candidate: &[u8]) -> Result<()> {
let current = object::File::parse(current).context("cannot inspect installed binary")?;
let candidate = object::File::parse(candidate).context("invalid downloaded binary")?;
ensure!(
candidate.kind() == object::ObjectKind::Executable
|| candidate.kind() == object::ObjectKind::Dynamic,
"download is not an executable"
);
ensure!(
candidate.format() == current.format()
&& candidate.architecture() == current.architecture()
&& candidate.is_64() == current.is_64()
&& candidate.is_little_endian() == current.is_little_endian(),
"downloaded binary has incompatible platform or architecture"
);
Ok(())
}
pub fn install(update: &Update, mut progress: impl FnMut(String)) -> Result<()> {
let exe = std::env::current_exe()?.canonicalize()?;
let parent = exe.parent().context("executable has no parent directory")?;
let lock = File::options()
.read(true)
.write(true)
.create(true)
.truncate(false)
.open(parent.join(".furumi-update.lock"))
.context("cannot write to installation directory")?;
fs2::FileExt::try_lock_exclusive(&lock).context("another furumi instance is updating")?;
// Keep staging on the destination filesystem and never overwrite a running image.
let staging = tempfile::Builder::new()
.prefix(".furumi-update-")
.tempdir_in(parent)
.context("cannot create update staging directory")?;
let client = client()?;
let sums = read_limited(
client
.get(&update.checksums.browser_download_url)
.send()?
.error_for_status()?,
1024 * 1024,
)?;
let expected = checksum(std::str::from_utf8(&sums)?, &update.asset.name)?;
let mut response = client
.get(&update.asset.browser_download_url)
.send()?
.error_for_status()?;
let total = response.content_length();
ensure!(
total.is_none_or(|size| size <= MAX_ARCHIVE),
"archive exceeds size limit"
);
let archive = staging.path().join("download");
let mut file = File::create(&archive)?;
let mut hash = Sha256::new();
let mut buffer = [0u8; 64 * 1024];
let mut downloaded = 0u64;
let mut reported = u64::MAX;
loop {
let count = response.read(&mut buffer)?;
if count == 0 {
break;
}
downloaded += count as u64;
ensure!(downloaded <= MAX_ARCHIVE, "archive exceeds size limit");
file.write_all(&buffer[..count])?;
hash.update(&buffer[..count]);
let mb = downloaded / (1024 * 1024);
if mb != reported {
progress(format!("Downloading: {mb} MiB"));
reported = mb;
}
}
file.sync_all()?;
drop(file);
ensure!(
format!("{:x}", hash.finalize()) == expected,
"SHA-256 mismatch; update was not installed"
);
progress("Verifying binary...".into());
let binary = staging.path().join(if cfg!(windows) {
"furumi.exe"
} else {
"furumi"
});
let mut output = File::create(&binary)?;
extract(&archive, &update.asset.name, &mut output)?;
drop(output);
validate_binary(&std::fs::read(&exe)?, &std::fs::read(&binary)?)?;
std::fs::set_permissions(&binary, std::fs::metadata(&exe)?.permissions())?;
progress("Installing...".into());
self_replace::self_replace(&binary).context("could not replace executable")?;
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn https_client_can_be_constructed() {
client().unwrap();
}
#[test]
fn extracts_only_expected_binary_from_release_archives() {
let temp = tempfile::tempdir().unwrap();
let binary = if cfg!(windows) {
"furumi.exe"
} else {
"furumi"
};
let expected = format!("furumi-windows-x86_64/{binary}");
let payload = b"test executable";
let zip_path = temp.path().join("test.zip");
let mut zip = zip::ZipWriter::new(File::create(&zip_path).unwrap());
zip.start_file("README.md", zip::write::SimpleFileOptions::default())
.unwrap();
zip.write_all(b"readme").unwrap();
zip.start_file(&expected, zip::write::SimpleFileOptions::default())
.unwrap();
zip.write_all(payload).unwrap();
zip.finish().unwrap();
let output = temp.path().join("output");
extract(
&zip_path,
"furumi-windows-x86_64-0.1.6.zip",
&mut File::create(&output).unwrap(),
)
.unwrap();
assert_eq!(std::fs::read(&output).unwrap(), payload);
assert!(
extract(
&zip_path,
"furumi-linux-x86_64-0.1.6.zip",
&mut File::create(&output).unwrap()
)
.is_err()
);
let tar_path = temp.path().join("test.tar.gz");
let gzip = flate2::write::GzEncoder::new(
File::create(&tar_path).unwrap(),
flate2::Compression::default(),
);
let mut tar = tar::Builder::new(gzip);
let mut header = tar::Header::new_gnu();
header.set_size(payload.len() as u64);
header.set_mode(0o755);
header.set_cksum();
tar.append_data(&mut header, &expected, payload.as_slice())
.unwrap();
tar.into_inner().unwrap().finish().unwrap();
extract(
&tar_path,
"furumi-windows-x86_64-0.1.6.tar.gz",
&mut File::create(&output).unwrap(),
)
.unwrap();
assert_eq!(std::fs::read(&output).unwrap(), payload);
}
#[test]
fn binary_validation_rejects_wrong_architecture() {
let current = std::fs::read(std::env::current_exe().unwrap()).unwrap();
let mut other = current.clone();
// Change only the machine field, preserving an otherwise valid executable.
if current.starts_with(b"MZ") {
let pe = u32::from_le_bytes(current[0x3c..0x40].try_into().unwrap()) as usize;
let machine: u16 = if current[pe + 4..pe + 6] == [0x64, 0x86] {
0xaa64
} else {
0x8664
};
other[pe + 4..pe + 6].copy_from_slice(&machine.to_le_bytes());
} else if current.starts_with(b"\x7fELF") {
let machine: u16 = if current[18..20] == [62, 0] { 183 } else { 62 };
other[18..20].copy_from_slice(&machine.to_le_bytes());
} else if current.starts_with(&[0xcf, 0xfa, 0xed, 0xfe]) {
let cpu: u32 = if current[4] == 7 {
0x0100000c
} else {
0x01000007
};
other[4..8].copy_from_slice(&cpu.to_le_bytes());
} else {
panic!("unsupported test executable format");
}
assert!(validate_binary(&current, &other).is_err());
}
#[test]
fn checksums_require_exact_unique_valid_entry() {
let hash = "ab".repeat(32);
assert_eq!(
checksum(&format!("{hash} app.zip\n"), "app.zip").unwrap(),
hash
);
assert!(checksum(&format!("{hash} other.zip"), "app.zip").is_err());
assert!(checksum("broken app.zip", "app.zip").is_err());
assert!(checksum(&format!("{hash} app.zip\n{hash} app.zip"), "app.zip").is_err());
}
#[test]
fn release_selection_respects_semver_and_assets() {
let release = |version: &str| Release {
tag_name: version.into(),
draft: false,
prerelease: false,
assets: vec![
Asset {
name: "furumi-linux-x86_64-0.1.10.tar.gz".into(),
browser_download_url: String::new(),
},
Asset {
name: "SHA256SUMS".into(),
browser_download_url: String::new(),
},
],
};
assert!(
select_release(release("v0.1.10"), "0.1.9", "linux", "x86_64")
.unwrap()
.is_some()
);
assert!(
select_release(release("v0.1.8"), "0.1.9", "linux", "x86_64")
.unwrap()
.is_none()
);
assert!(
select_release(release("v0.2.0-beta.1"), "0.1.9", "linux", "x86_64")
.unwrap()
.is_none()
);
assert!(select_release(release("v0.1.10"), "0.1.9", "linux", "aarch64").is_err());
}
#[test]
fn binary_validation_rejects_corrupt_download() {
let current = std::fs::read(std::env::current_exe().unwrap()).unwrap();
validate_binary(&current, &current).unwrap();
assert!(validate_binary(&current, b"not a binary").is_err());
}
}