Refactor to split client from server

This commit is contained in:
Alexander
2026-07-01 16:31:55 +02:00
parent d35aa2aecb
commit 8dba0c097d
56 changed files with 441 additions and 326 deletions
@@ -0,0 +1,102 @@
use std::{net::SocketAddr, path::PathBuf};
use anyhow::{Context, Result};
use clap::Parser;
use musicfs_core::logging::{LogConfig, init};
use musicfs_server::server::{
state::ServerState,
transport::{self, TransportArgs},
watcher::ServerWatcher,
};
use tokio::signal::unix::{SignalKind, signal};
use tracing::{error, info};
#[derive(Parser, Debug)]
#[command(version, about, long_about = None)]
struct Args {
/// Directory containing the music library to serve.
#[arg(short, long, required = true)]
source: PathBuf,
/// Address:port the transport should listen on (e.g. 0.0.0.0:50051).
#[arg(short, long, default_value = "0.0.0.0:50051")]
listen: SocketAddr,
/// Name of the transport implementation to use. Today only "grpc".
#[arg(short, long, default_value = "grpc")]
transport: String,
/// Directory for daily-rotated log files.
#[arg(long, default_value = "./logs")]
log_dir: PathBuf,
}
#[tokio::main]
async fn main() -> Result<()> {
let args = Args::parse();
// Bind the guard for the whole process so the non-blocking file writer
// flushes on exit. Initialized before anything else so even early
// failures land in the log.
let _guard = init(LogConfig {
log_dir: args.log_dir.clone(),
file_prefix: "musicfs-server".to_string(),
max_files: 7,
});
let source = args.source.clone();
info!(
source = %source.display(),
transport = %args.transport,
listen = %args.listen,
log_dir = %args.log_dir.display(),
"musicfs-server starting"
);
if !source.is_dir() {
error!(source = %source.display(), "source is not a readable directory");
return Err(anyhow::anyhow!(
"source is not a readable directory: {}",
source.display()
));
}
let state = ServerState::new();
state
.replace_all(&source)
.with_context(|| format!("initial scan of {}", source.display()))?;
info!(entries = state.manifest().len(), "initial scan complete");
let (_watcher, events_rx) = ServerWatcher::spawn(source.clone(), state.clone());
let transport = transport::build(
&args.transport,
TransportArgs {
listen: args.listen,
},
)
.with_context(|| format!("building transport {:?}", args.transport))?;
info!(
transport = %args.transport,
listen = %args.listen,
"transport ready"
);
let server_task = tokio::spawn(async move {
if let Err(e) = transport.run(state, events_rx).await {
error!(error = %e, "transport ended with error");
}
});
let mut sigint = signal(SignalKind::interrupt()).expect("register SIGINT");
let mut sigterm = signal(SignalKind::terminate()).expect("register SIGTERM");
tokio::select! {
_ = sigint.recv() => info!("received SIGINT, shutting down"),
_ = sigterm.recv() => info!("received SIGTERM, shutting down"),
_ = server_task => info!("transport task exited"),
}
return Ok(());
}
+1
View File
@@ -0,0 +1 @@
pub mod server;
@@ -0,0 +1,93 @@
use musicfs_core::compute_item_hash;
use musicfs_core::music::metadata::MusicMetadata;
/// One row of the in-memory manifest: every field the client needs to
/// reconstruct an `Item` whose hash matches the server's hash.
///
/// Field semantics:
/// - `id` is the server filesystem inode. The client uses it as the FUSE
/// inode, which keeps `compute_hash` inputs identical on both sides.
/// - `rel_path` is the path relative to the server's `--source` root, using
/// `/` as separator. The client stores it as `Item.original_path` and uses
/// it as the byte-source locator.
/// - `mtime` / `ctime` / `crtime` are seconds since the Unix epoch. These are
/// the exact three time fields `compute_hash` consumes.
/// - `size` is the real on-disk file size in bytes (the client's virtual size
/// is derived from `music_metadata.virtual_size(size)` when present).
/// - `music_metadata` is the fully-encoded metadata produced server-side via
/// `parse_music_metadata_for_path`. The client never parses audio.
///
/// Wire format conversion lives in `transport/grpc.rs` (`From<ManifestEntry>`
/// for the generated proto type).
#[derive(Debug, Clone)]
pub struct ManifestEntry {
pub id: u64,
pub rel_path: String,
pub size: u64,
pub mtime: u64,
pub ctime: u64,
pub crtime: u64,
pub music_metadata: Option<MusicMetadata>,
}
impl ManifestEntry {
pub fn hash(&self) -> u64 {
return compute_item_hash(
self.id,
self.rel_path.as_bytes(),
self.ctime,
self.mtime,
self.crtime,
);
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn manifest_entry_clones_fields_into_copy() {
let entry = ManifestEntry {
id: 42,
rel_path: "Artist/Album/track.flac".to_string(),
size: 12345,
mtime: 1_700_000_000,
ctime: 1_699_999_000,
crtime: 1_699_990_000,
music_metadata: None,
};
let cloned = entry.clone();
assert_eq!(cloned.id, entry.id);
assert_eq!(cloned.rel_path, entry.rel_path);
assert_eq!(cloned.size, entry.size);
assert_eq!(cloned.mtime, entry.mtime);
assert_eq!(cloned.ctime, entry.ctime);
assert_eq!(cloned.crtime, entry.crtime);
assert!(cloned.music_metadata.is_none());
}
#[test]
fn manifest_entry_carries_music_metadata_when_present() {
let mm = MusicMetadata {
artist: vec!["Test Artist".to_string()],
album: "Test Album".to_string(),
track_title: "Test Title".to_string(),
track_number: 3,
..MusicMetadata::default()
};
let entry = ManifestEntry {
id: 1,
rel_path: "a/b.flac".to_string(),
size: 100,
mtime: 0,
ctime: 0,
crtime: 0,
music_metadata: Some(mm.clone()),
};
assert_eq!(entry.music_metadata.as_ref().unwrap().album, "Test Album");
assert_eq!(entry.music_metadata.as_ref().unwrap().track_number, 3);
}
}
+5
View File
@@ -0,0 +1,5 @@
pub mod manifest;
pub mod range;
pub mod state;
pub mod transport;
pub mod watcher;
+175
View File
@@ -0,0 +1,175 @@
/// Parsed `Range: bytes=...` request header.
///
/// Supported forms (RFC 7233):
/// - `bytes=a-b` → `StartEnd(a, b)` (inclusive)
/// - `bytes=a-` → `Start(a)`
/// - `bytes=-N` → `Suffix(N)` (last N bytes)
///
/// Returns `None` if the header is missing, uses a unit other than `bytes`,
/// or fails to parse. Multiple ranges (`bytes=a-b,c-d`) are not supported —
/// the first is used and the rest ignored.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ByteRange {
StartEnd(u64, u64),
Start(u64),
Suffix(u64),
}
pub fn parse_range_header(header: &str) -> Option<ByteRange> {
let header = header.trim();
let rest = header.strip_prefix("bytes=")?;
let first = rest.split(',').next()?.trim();
let (left, right) = first.split_once('-')?;
let left = left.trim();
let right = right.trim();
return match (left.is_empty(), right.is_empty()) {
(false, false) => Some(ByteRange::StartEnd(left.parse().ok()?, right.parse().ok()?)),
(false, true) => Some(ByteRange::Start(left.parse().ok()?)),
(true, false) => Some(ByteRange::Suffix(right.parse().ok()?)),
(true, true) => None,
};
}
/// Resolve a parsed range against a real file size, yielding an absolute
/// `(start, length)` pair suitable for `seek + read`.
///
/// Returns `None` if the resolved range is unsatisfiable (e.g. start past
/// end of file). The returned `start` is clamped to `[0, size]` and `length`
/// is clamped to not exceed `size - start`.
pub fn resolve_range(range: ByteRange, size: u64) -> Option<(u64, u64)> {
let (start, end_inclusive) = match range {
ByteRange::StartEnd(a, b) => {
if a >= size || a > b {
return None;
}
(a, b.min(size - 1))
}
ByteRange::Start(a) => {
if a >= size {
return None;
}
(a, size - 1)
}
ByteRange::Suffix(n) => {
if n == 0 || size == 0 {
return None;
}
// A suffix larger than the file clamps to the whole file.
let start = size.saturating_sub(n);
(start, size - 1)
}
};
return Some((start, end_inclusive - start + 1));
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn parse_start_end() {
assert_eq!(
parse_range_header("bytes=0-499"),
Some(ByteRange::StartEnd(0, 499))
);
assert_eq!(
parse_range_header("bytes=500-999"),
Some(ByteRange::StartEnd(500, 999))
);
}
#[test]
fn parse_start_only() {
assert_eq!(
parse_range_header("bytes=9500-"),
Some(ByteRange::Start(9500))
);
}
#[test]
fn parse_suffix() {
assert_eq!(
parse_range_header("bytes=-500"),
Some(ByteRange::Suffix(500))
);
}
#[test]
fn parse_ignores_additional_ranges_after_first() {
assert_eq!(
parse_range_header("bytes=0-499,1000-1499"),
Some(ByteRange::StartEnd(0, 499))
);
}
#[test]
fn parse_returns_none_for_non_bytes_unit() {
assert_eq!(parse_range_header("items=0-4"), None);
}
#[test]
fn parse_returns_none_for_malformed() {
assert_eq!(parse_range_header("not a range"), None);
assert_eq!(parse_range_header("bytes="), None);
assert_eq!(parse_range_header("bytes=-"), None);
}
#[test]
fn resolve_start_end_in_bounds() {
assert_eq!(
resolve_range(ByteRange::StartEnd(0, 499), 1000),
Some((0, 500))
);
assert_eq!(
resolve_range(ByteRange::StartEnd(100, 199), 1000),
Some((100, 100))
);
}
#[test]
fn resolve_start_end_clamps_end_to_size() {
assert_eq!(
resolve_range(ByteRange::StartEnd(900, 2000), 1000),
Some((900, 100))
);
}
#[test]
fn resolve_start_end_unsatisfiable_when_start_past_size() {
assert_eq!(resolve_range(ByteRange::StartEnd(1500, 2000), 1000), None);
}
#[test]
fn resolve_start_only() {
assert_eq!(resolve_range(ByteRange::Start(900), 1000), Some((900, 100)));
assert_eq!(resolve_range(ByteRange::Start(0), 1000), Some((0, 1000)));
}
#[test]
fn resolve_start_unsatisfiable_when_past_size() {
assert_eq!(resolve_range(ByteRange::Start(1500), 1000), None);
}
#[test]
fn resolve_suffix() {
assert_eq!(
resolve_range(ByteRange::Suffix(500), 1000),
Some((500, 500))
);
assert_eq!(resolve_range(ByteRange::Suffix(1), 1000), Some((999, 1)));
}
#[test]
fn resolve_suffix_larger_than_size_clamps_to_full_file() {
assert_eq!(
resolve_range(ByteRange::Suffix(2000), 1000),
Some((0, 1000))
);
}
#[test]
fn resolve_suffix_zero_unsatisfiable() {
assert_eq!(resolve_range(ByteRange::Suffix(0), 1000), None);
}
}
+260
View File
@@ -0,0 +1,260 @@
use std::{
collections::HashMap,
fs, io,
path::{Path, PathBuf},
sync::{Arc, Mutex},
time::{SystemTime, UNIX_EPOCH},
};
use std::os::unix::fs::MetadataExt;
use crate::server::manifest::ManifestEntry;
use musicfs_core::FileAttrs;
use musicfs_core::music::parse::parse_music_metadata_for_path;
use tracing::info;
/// Server-side entry for one file. Built once on startup from a directory
/// scan and refreshed by the watcher on inotify events.
#[derive(Debug, Clone, PartialEq)]
pub struct FileEntry {
pub abs_path: PathBuf,
pub rel_path: String,
pub attrs: FileAttrs,
pub music_metadata: Option<musicfs_core::music::metadata::MusicMetadata>,
}
#[derive(Clone)]
pub struct ServerState {
inner: Arc<Mutex<HashMap<u64, FileEntry>>>,
}
impl ServerState {
pub fn new() -> Self {
return ServerState {
inner: Arc::new(Mutex::new(HashMap::new())),
};
}
/// Replace the entire map with a fresh scan of `source`. Used on startup
/// and on watcher-driven reconciliations.
pub fn replace_all(&self, source: &Path) -> io::Result<bool> {
let entries = scan_directory(source)?;
let new_map: HashMap<u64, FileEntry> = entries.into_iter().collect();
let mut map = self.inner.lock().unwrap();
if *map == new_map {
return Ok(false);
}
let count = new_map.len();
*map = new_map;
info!(count, "server state: scan complete (changed)");
return Ok(true);
}
/// Snapshot the current state into a manifest. Order is by inode ascending
/// so two scans of an unchanged library serialize identically.
pub fn manifest(&self) -> Vec<ManifestEntry> {
let map = self.inner.lock().unwrap();
let mut entries: Vec<ManifestEntry> = map
.iter()
.map(|(id, file)| ManifestEntry {
id: *id,
rel_path: file.rel_path.clone(),
size: file.attrs.size,
mtime: secs(file.attrs.mtime),
ctime: secs(file.attrs.ctime),
crtime: secs(file.attrs.crtime),
music_metadata: file.music_metadata.clone(),
})
.collect();
entries.sort_by_key(|e| e.id);
return entries;
}
pub fn lookup(&self, id: u64) -> Option<FileEntry> {
return self.inner.lock().unwrap().get(&id).cloned();
}
}
fn secs(time: SystemTime) -> u64 {
return time
.duration_since(UNIX_EPOCH)
.map(|d| d.as_secs())
.unwrap_or(0);
}
fn scan_directory(source: &Path) -> io::Result<Vec<(u64, FileEntry)>> {
let mut out = Vec::new();
walk(source, source, &mut out)?;
return Ok(out);
}
fn walk(source: &Path, dir: &Path, out: &mut Vec<(u64, FileEntry)>) -> io::Result<()> {
for entry in fs::read_dir(dir)? {
let entry = entry?;
let path = entry.path();
let file_type = entry.file_type()?;
if file_type.is_dir() {
walk(source, &path, out)?;
continue;
}
if !file_type.is_file() {
continue;
}
let metadata = entry.metadata()?;
let rel_path = path
.strip_prefix(source)
.map(|p| p.to_path_buf())
.unwrap_or_else(|_| path.clone())
.to_string_lossy()
.replace('\\', "/");
let inode = metadata.ino();
let music_metadata = parse_music_metadata_for_path(&path);
let file_entry = FileEntry {
abs_path: path,
rel_path,
attrs: FileAttrs::from(&metadata),
music_metadata,
};
out.push((inode, file_entry));
}
return Ok(());
}
#[cfg(test)]
mod tests {
use super::*;
use std::io::Write;
#[test]
fn replace_all_loads_files_into_state() {
let tmp = tempfile::tempdir().unwrap();
let source = tmp.path();
let mut f = fs::File::create(source.join("a.txt")).unwrap();
f.write_all(b"hello").unwrap();
drop(f);
let mut f = fs::File::create(source.join("b.txt")).unwrap();
f.write_all(b"world!").unwrap();
drop(f);
let state = ServerState::new();
state.replace_all(source).unwrap();
let manifest = state.manifest();
assert_eq!(manifest.len(), 2);
assert!(manifest.iter().any(|e| e.rel_path == "a.txt"));
assert!(manifest.iter().any(|e| e.rel_path == "b.txt"));
}
#[test]
fn manifest_entries_sorted_by_id_ascending() {
let tmp = tempfile::tempdir().unwrap();
let source = tmp.path();
// Create files in any order; inode order is determined by the FS.
for name in ["z.txt", "a.txt", "m.txt"] {
fs::write(source.join(name), b"x").unwrap();
}
let state = ServerState::new();
state.replace_all(source).unwrap();
let ids: Vec<u64> = state.manifest().iter().map(|e| e.id).collect();
let mut sorted = ids.clone();
sorted.sort();
assert_eq!(ids, sorted);
}
#[test]
fn replace_all_clears_existing_entries() {
let tmp = tempfile::tempdir().unwrap();
let source = tmp.path();
fs::write(source.join("a.txt"), b"x").unwrap();
let state = ServerState::new();
state.replace_all(source).unwrap();
assert_eq!(state.manifest().len(), 1);
fs::remove_file(source.join("a.txt")).unwrap();
state.replace_all(source).unwrap();
assert_eq!(state.manifest().len(), 0);
}
#[test]
fn replace_all_returns_true_on_first_scan() {
let tmp = tempfile::tempdir().unwrap();
let source = tmp.path();
fs::write(source.join("a.txt"), b"x").unwrap();
let state = ServerState::new();
assert!(state.replace_all(source).unwrap());
}
#[test]
fn replace_all_returns_false_when_unchanged() {
let tmp = tempfile::tempdir().unwrap();
let source = tmp.path();
fs::write(source.join("a.txt"), b"hello").unwrap();
fs::write(source.join("b.txt"), b"world").unwrap();
let state = ServerState::new();
state.replace_all(source).unwrap();
assert!(!state.replace_all(source).unwrap());
}
#[test]
fn replace_all_returns_true_after_file_added() {
let tmp = tempfile::tempdir().unwrap();
let source = tmp.path();
fs::write(source.join("a.txt"), b"x").unwrap();
let state = ServerState::new();
state.replace_all(source).unwrap();
fs::write(source.join("b.txt"), b"y").unwrap();
assert!(state.replace_all(source).unwrap());
}
#[test]
fn replace_all_returns_true_after_file_removed() {
let tmp = tempfile::tempdir().unwrap();
let source = tmp.path();
fs::write(source.join("a.txt"), b"x").unwrap();
fs::write(source.join("b.txt"), b"y").unwrap();
let state = ServerState::new();
state.replace_all(source).unwrap();
fs::remove_file(source.join("a.txt")).unwrap();
assert!(state.replace_all(source).unwrap());
}
#[test]
fn replace_all_returns_true_after_file_renamed() {
let tmp = tempfile::tempdir().unwrap();
let source = tmp.path();
fs::write(source.join("old.txt"), b"x").unwrap();
let state = ServerState::new();
state.replace_all(source).unwrap();
assert!(!state.replace_all(source).unwrap());
fs::rename(source.join("old.txt"), source.join("new.txt")).unwrap();
assert!(state.replace_all(source).unwrap());
}
#[test]
fn replace_all_returns_true_after_content_modified() {
let tmp = tempfile::tempdir().unwrap();
let source = tmp.path();
fs::write(source.join("a.txt"), b"original").unwrap();
let state = ServerState::new();
state.replace_all(source).unwrap();
assert!(!state.replace_all(source).unwrap());
fs::write(source.join("a.txt"), b"modified content").unwrap();
assert!(state.replace_all(source).unwrap());
}
}
@@ -0,0 +1,328 @@
use std::{
fs,
io::{Read, Seek, SeekFrom},
path::PathBuf,
pin::Pin,
};
use anyhow::Result;
use async_trait::async_trait;
use tokio::sync::broadcast::Receiver;
use tonic::{Request, Response, Status, transport::Server};
use crate::server::manifest::ManifestEntry as DomainManifestEntry;
use crate::server::state::ServerState;
use crate::server::transport::{MusicTransport, TransportArgs};
use crate::server::watcher::{ChangeEvent, ChangeKind};
use musicfs_core::music::metadata::MusicMetadata;
use musicfs_proto as proto_types;
use musicfs_proto::MusicFs as MusicFsTrait;
use musicfs_proto::{
ChangeEvent as ProtoChangeEvent, FileChunk, GetFileRequest, GetManifestRequest,
GetMetadataRequest, InodeHash, ManifestEntry, MusicFsServer,
MusicMetadata as ProtoMusicMetadata, PictureDataRange, ReconcileRequest, ReconcileResponse,
SubscribeEventsRequest,
};
use tracing::{debug, info, warn};
pub struct GrpcTransport {
args: TransportArgs,
}
impl GrpcTransport {
pub fn new(args: TransportArgs) -> Self {
return GrpcTransport { args };
}
}
#[async_trait]
impl MusicTransport for GrpcTransport {
async fn run(self: Box<Self>, state: ServerState, events: Receiver<ChangeEvent>) -> Result<()> {
let service = MusicFsService {
state,
events: tokio::sync::Mutex::new(events),
};
let listen = self.args.listen;
info!(%listen, "musicfs-server gRPC listening");
// Standard gRPC Health Checking Protocol (grpc.health.v1.Health).
let (reporter, health_service) = tonic_health::server::health_reporter();
reporter
.set_serving::<MusicFsServer<MusicFsService>>()
.await;
// v1alpha covers older clients; v1 is the current standard.
let reflection_v1 = tonic_reflection::server::Builder::configure()
.register_encoded_file_descriptor_set(musicfs_proto::musicfs::FILE_DESCRIPTOR_SET)
.register_encoded_file_descriptor_set(tonic_health::pb::FILE_DESCRIPTOR_SET)
.build_v1()?;
let reflection_v1alpha = tonic_reflection::server::Builder::configure()
.register_encoded_file_descriptor_set(musicfs_proto::musicfs::FILE_DESCRIPTOR_SET)
.register_encoded_file_descriptor_set(tonic_health::pb::FILE_DESCRIPTOR_SET)
.build_v1alpha()?;
Server::builder()
.add_service(health_service)
.add_service(reflection_v1)
.add_service(reflection_v1alpha)
.add_service(MusicFsServer::new(service))
.serve(listen)
.await?;
return Ok(());
}
}
struct MusicFsService {
state: ServerState,
events: tokio::sync::Mutex<Receiver<ChangeEvent>>,
}
type BoxStream<T> = Pin<Box<dyn tokio_stream::Stream<Item = Result<T, Status>> + Send>>;
#[tonic::async_trait]
impl MusicFsTrait for MusicFsService {
type GetManifestStream = BoxStream<ManifestEntry>;
type GetMetadataStream = BoxStream<ManifestEntry>;
type GetFileStream = BoxStream<FileChunk>;
type SubscribeEventsStream = BoxStream<ProtoChangeEvent>;
async fn reconcile(
&self,
request: Request<ReconcileRequest>,
) -> Result<Response<ReconcileResponse>, Status> {
let client_entries: std::collections::HashMap<u64, u64> = request
.into_inner()
.entries
.into_iter()
.map(|ih| (ih.inode, ih.hash))
.collect();
debug!(client_entries = client_entries.len(), "reconcile");
let server_manifest = self.state.manifest();
let mut changed: Vec<InodeHash> = Vec::new();
let mut deleted: Vec<u64> = Vec::new();
for entry in &server_manifest {
let entry_hash = entry.hash();
match client_entries.get(&entry.id) {
Some(client_hash) if *client_hash == entry_hash => {}
_ => {
changed.push(InodeHash {
inode: entry.id,
hash: entry_hash,
});
}
}
}
for client_inode in client_entries.keys() {
if !server_manifest.iter().any(|e| &e.id == client_inode) {
deleted.push(*client_inode);
}
}
debug!(
changed = changed.len(),
deleted = deleted.len(),
"reconcile result"
);
return Ok(Response::new(ReconcileResponse { changed, deleted }));
}
async fn get_metadata(
&self,
request: Request<GetMetadataRequest>,
) -> Result<Response<Self::GetMetadataStream>, Status> {
let wanted: std::collections::HashSet<u64> =
request.into_inner().inodes.into_iter().collect();
debug!(wanted = wanted.len(), "get_metadata");
let entries: Vec<DomainManifestEntry> = self
.state
.manifest()
.into_iter()
.filter(|e| wanted.contains(&e.id))
.collect();
let stream = tokio_stream::iter(
entries
.into_iter()
.map(|e| Ok::<ManifestEntry, Status>(e.into())),
);
return Ok(Response::new(Box::pin(stream)));
}
async fn get_manifest(
&self,
_request: Request<GetManifestRequest>,
) -> Result<Response<Self::GetManifestStream>, Status> {
debug!("get_manifest");
let entries: Vec<DomainManifestEntry> = self.state.manifest();
let stream = tokio_stream::iter(
entries
.into_iter()
.map(|e| Ok::<ManifestEntry, Status>(e.into())),
);
return Ok(Response::new(Box::pin(stream)));
}
async fn get_file(
&self,
request: Request<GetFileRequest>,
) -> Result<Response<Self::GetFileStream>, Status> {
let req = request.into_inner();
let id = req.id;
debug!(%id, "get_file");
let entry = match self.state.lookup(id) {
Some(e) => e,
None => {
debug!(%id, "get_file: not found");
return Err(Status::not_found(format!("no file with id {id}")));
}
};
let abs_path: PathBuf = entry.abs_path.clone();
let total_size = entry.attrs.size;
let start = if req.start == 0 && req.length == 0 {
0u64
} else {
req.start
};
let length = if req.start == 0 && req.length == 0 {
total_size
} else if req.length == 0 {
total_size.saturating_sub(start)
} else {
req.length
};
let chunk =
tokio::task::spawn_blocking(move || read_range_blocking(&abs_path, start, length))
.await
.map_err(|e| {
warn!(%id, error = %e, "get_file: join blocking read failed");
Status::internal(format!("join blocking read: {e}"))
})?
.map_err(|e| {
warn!(%id, error = %e, "get_file: read range failed");
Status::internal(format!("read file range: {e}"))
})?;
// Split into <=1MB messages: gRPC's default max is 4MB and music files
// routinely exceed it.
const CHUNK_SIZE: usize = 1024 * 1024;
let chunk_stream = tokio_stream::iter(
chunk
.chunks(CHUNK_SIZE)
.enumerate()
.map(|(i, c)| {
Ok::<FileChunk, Status>(FileChunk {
data: c.to_vec(),
offset: start + (i * CHUNK_SIZE) as u64,
total_size,
})
})
.collect::<Vec<_>>(),
);
return Ok(Response::new(Box::pin(chunk_stream)));
}
async fn subscribe_events(
&self,
_request: Request<SubscribeEventsRequest>,
) -> Result<Response<Self::SubscribeEventsStream>, Status> {
debug!("subscribe_events");
// Take an owned receiver (the guard is dropped at the end of this
// statement) so it can move into the 'static spawned task, and so each
// subscriber gets its own receiver rather than contending on one.
let mut rx = self.events.lock().await.resubscribe();
let (tx, rx_stream) =
tokio::sync::mpsc::unbounded_channel::<Result<ProtoChangeEvent, Status>>();
tokio::spawn(async move {
loop {
match rx.recv().await {
Ok(event) => {
let proto = ProtoChangeEvent {
kind: kind_to_string(event.kind),
};
if tx.send(Ok(proto)).is_err() {
break;
}
}
Err(tokio::sync::broadcast::error::RecvError::Lagged(n)) => {
warn!(lagged = n, "subscribe_events: lagged; continuing");
continue;
}
Err(tokio::sync::broadcast::error::RecvError::Closed) => {
debug!("subscribe_events: broadcast closed; ending stream");
break;
}
}
}
});
let stream = tokio_stream::wrappers::UnboundedReceiverStream::new(rx_stream);
return Ok(Response::new(Box::pin(stream)));
}
}
fn read_range_blocking(
path: &std::path::Path,
start: u64,
length: u64,
) -> std::io::Result<Vec<u8>> {
let mut file = fs::File::open(path)?;
file.seek(SeekFrom::Start(start))?;
let mut buf = vec![0u8; length as usize];
let mut filled = 0usize;
while filled < buf.len() {
let n = file.read(&mut buf[filled..])?;
if n == 0 {
buf.truncate(filled);
break;
}
filled += n;
}
return Ok(buf);
}
fn kind_to_string(kind: ChangeKind) -> String {
return match kind {
ChangeKind::Create => "create".to_string(),
ChangeKind::Modify => "modify".to_string(),
ChangeKind::Remove => "remove".to_string(),
};
}
impl From<DomainManifestEntry> for ManifestEntry {
fn from(entry: DomainManifestEntry) -> Self {
return ManifestEntry {
id: entry.id,
rel_path: entry.rel_path,
size: entry.size,
mtime: entry.mtime,
ctime: entry.ctime,
crtime: entry.crtime,
music_metadata: entry.music_metadata.map(music_metadata_to_proto),
};
}
}
fn music_metadata_to_proto(mm: MusicMetadata) -> ProtoMusicMetadata {
ProtoMusicMetadata {
artist: mm.artist,
album_artist: mm.album_artist,
album: mm.album,
track_number: mm.track_number,
track_title: mm.track_title,
other_tags: mm.other_tags,
header: mm.header,
picture_block_headers: mm.picture_block_headers,
picture_data_ranges: mm
.picture_data_ranges
.into_iter()
.map(|(offset, length)| PictureDataRange { offset, length })
.collect(),
real_audio_start: mm.real_audio_start,
vorbis_comment_offset: mm.vorbis_comment_offset,
vorbis_comment_length: mm.vorbis_comment_length,
}
}
#[allow(unused_imports)]
use proto_types as _proto_types_anchor;
@@ -0,0 +1,39 @@
pub mod grpc;
use std::net::SocketAddr;
use anyhow::{Result, anyhow};
use async_trait::async_trait;
use tokio::sync::broadcast::Receiver;
use crate::server::state::ServerState;
use crate::server::watcher::ChangeEvent;
/// Transport-agnostic musicfs server. A single implementation is wired in
/// today (`grpc`); the factory in [`build`] is the seam where additional
/// transports (HTTP/3, raw QUIC, in-process for tests, ...) get plugged in
/// without touching the binary.
///
/// `run` consumes `self` — a transport serves exactly once. It receives the
/// shared [`ServerState`] and a fresh broadcast receiver for change events.
#[async_trait]
pub trait MusicTransport: Send + Sync + 'static {
async fn run(self: Box<Self>, state: ServerState, events: Receiver<ChangeEvent>) -> Result<()>;
}
/// Arguments every transport understands. Transports may extend this with
/// their own configuration via their constructors.
#[derive(Debug, Clone, Copy)]
pub struct TransportArgs {
pub listen: SocketAddr,
}
/// Construct the named transport. The list of accepted names is intentionally
/// discoverable (a single match arm) so adding a transport means adding an
/// arm here plus a module under `transport/`.
pub fn build(name: &str, args: TransportArgs) -> Result<Box<dyn MusicTransport>> {
return match name {
"grpc" => Ok(Box::new(grpc::GrpcTransport::new(args))),
other => Err(anyhow!("unknown transport: {other}")),
};
}
+119
View File
@@ -0,0 +1,119 @@
use std::{
path::PathBuf,
sync::{Arc, Mutex},
thread,
time::Duration,
};
use notify::{EventKind, RecursiveMode, Watcher};
use tokio::sync::broadcast;
use crate::server::state::ServerState;
use tracing::{error, info};
const POLL_INTERVAL: Duration = Duration::from_secs(15);
/// A change observed by the watcher. Pushed onto the broadcast channel for
/// `/events` subscribers. The client treats these as wake-ups: correctness
/// always rests on the subsequent `/manifest` hash diff.
#[derive(Debug, Clone)]
pub struct ChangeEvent {
pub kind: ChangeKind,
}
#[derive(Debug, Clone, Copy)]
pub enum ChangeKind {
Create,
Modify,
Remove,
}
/// Server-side file watcher. Owns a background std thread driving `notify`
/// (inotify on Linux). On any filesystem event under `source` it:
/// 1. Rebuilds the shared [`ServerState`] from a fresh scan.
/// 2. Broadcasts a [`ChangeEvent`] so connected `/events` clients wake up.
///
/// The state rebuild is full-scan rather than incremental. This matches the
/// existing `LocalOriginFileWatcher` pattern, keeps the watcher simple, and
/// is cheap on a LAN-scale library. Hash-diff reconciliation on the client
/// absorbs any over-reporting.
pub struct ServerWatcher {
_events_tx: broadcast::Sender<ChangeEvent>,
_worker: Arc<Mutex<Option<thread::JoinHandle<()>>>>,
}
impl ServerWatcher {
/// Spawn the watcher. Returns the broadcast receiver that `/events`
/// handlers subscribe to.
pub fn spawn(source: PathBuf, state: ServerState) -> (Self, broadcast::Receiver<ChangeEvent>) {
let (events_tx, events_rx) = broadcast::channel(64);
let worker_tx = events_tx.clone();
let handle = thread::spawn(move || {
run_watcher_loop(source, state, worker_tx);
});
let watcher = ServerWatcher {
_events_tx: events_tx,
_worker: Arc::new(Mutex::new(Some(handle))),
};
return (watcher, events_rx);
}
}
fn run_watcher_loop(
source: PathBuf,
state: ServerState,
events_tx: broadcast::Sender<ChangeEvent>,
) {
let (tx, rx) = std::sync::mpsc::channel();
let mut watcher = match notify::recommended_watcher(tx) {
Ok(w) => w,
Err(e) => {
error!(error = %e, "server watcher: failed to create inotify watcher");
return;
}
};
if let Err(e) = watcher.watch(&source, RecursiveMode::Recursive) {
error!(source = %source.display(), error = %e, "server watcher: failed to watch source");
return;
}
info!(source = %source.display(), "server watcher: scanning for changes");
loop {
match rx.recv_timeout(POLL_INTERVAL) {
Ok(res) => match res {
Ok(event) => match event.kind {
EventKind::Create(_) | EventKind::Modify(_) | EventKind::Remove(_) => {
let kind = match event.kind {
EventKind::Create(_) => ChangeKind::Create,
EventKind::Modify(_) => ChangeKind::Modify,
EventKind::Remove(_) => ChangeKind::Remove,
_ => continue,
};
match state.replace_all(&source) {
Ok(true) => {
let _ = events_tx.send(ChangeEvent { kind });
}
Ok(false) => {}
Err(e) => {
error!(error = %e, "server watcher: state refresh failed");
}
}
}
_ => {}
},
Err(e) => error!(error = %e, "server watcher: inotify error"),
},
Err(std::sync::mpsc::RecvTimeoutError::Timeout) => match state.replace_all(&source) {
Ok(true) => {
let _ = events_tx.send(ChangeEvent {
kind: ChangeKind::Modify,
});
}
Ok(false) => {}
Err(e) => error!(error = %e, "server watcher: poll rescan failed"),
},
Err(std::sync::mpsc::RecvTimeoutError::Disconnected) => break,
}
}
}