Bump protocols. Added proto status
This commit is contained in:
@@ -18,6 +18,8 @@ use crate::library::Library;
|
||||
|
||||
/// ALPN of the audio streaming protocol (shared with furumi-fd).
|
||||
pub const AUDIO_ALPN: &[u8] = b"furumi-fd/audio/1";
|
||||
/// Version of the audio transfer stream protocol.
|
||||
pub const AUDIO_PROTOCOL_VERSION: u16 = 1;
|
||||
|
||||
/// Maximum size of a JSON protocol line (request or response header).
|
||||
const MAX_PROTOCOL_LINE: usize = 4096;
|
||||
|
||||
@@ -0,0 +1,190 @@
|
||||
//! Informational publication and observation of protocol versions.
|
||||
|
||||
use std::collections::BTreeMap;
|
||||
use std::sync::{Arc, Mutex, MutexGuard};
|
||||
use std::time::Duration;
|
||||
|
||||
use anyhow::{Context, Result};
|
||||
pub use music_dht::capabilities::CAPABILITIES_ALPN;
|
||||
use music_dht::capabilities::{
|
||||
CAPABILITIES_PROTOCOL_VERSION, CapabilityManifest, CapabilityMessage, read_message,
|
||||
write_message,
|
||||
};
|
||||
use music_dht::{ByteStream, EndpointId, MusicDhtService, StreamAcceptor};
|
||||
|
||||
const PROBE_INTERVAL: Duration = Duration::from_secs(30);
|
||||
|
||||
#[derive(Debug, Clone, Default, PartialEq, Eq)]
|
||||
pub struct ProtocolVersions {
|
||||
pub local: BTreeMap<String, u16>,
|
||||
pub observed: BTreeMap<String, u16>,
|
||||
pub observed_peers: usize,
|
||||
pub newer: Vec<NewerProtocol>,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, PartialEq, Eq)]
|
||||
pub struct NewerProtocol {
|
||||
pub id: String,
|
||||
pub local: u16,
|
||||
pub observed: u16,
|
||||
}
|
||||
|
||||
impl ProtocolVersions {
|
||||
pub fn snapshot(observed: &ObservedVersions) -> Self {
|
||||
let local = local_manifest().protocols;
|
||||
let observed_versions = lock(&observed.versions).clone();
|
||||
let newer = observed_versions
|
||||
.iter()
|
||||
.filter_map(|(id, remote)| {
|
||||
let local_version = local.get(id)?;
|
||||
(*remote > *local_version).then(|| NewerProtocol {
|
||||
id: id.clone(),
|
||||
local: *local_version,
|
||||
observed: *remote,
|
||||
})
|
||||
})
|
||||
.collect();
|
||||
Self {
|
||||
local,
|
||||
observed: observed_versions,
|
||||
observed_peers: lock(&observed.peers).len(),
|
||||
newer,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Default)]
|
||||
pub struct ObservedVersions {
|
||||
versions: Mutex<BTreeMap<String, u16>>,
|
||||
peers: Mutex<BTreeMap<String, String>>,
|
||||
}
|
||||
|
||||
fn local_manifest() -> CapabilityManifest {
|
||||
CapabilityManifest::frid("furumi", env!("CARGO_PKG_VERSION"))
|
||||
.with_protocol("audio", super::audio::AUDIO_PROTOCOL_VERSION)
|
||||
}
|
||||
|
||||
pub async fn serve(mut acceptor: StreamAcceptor) {
|
||||
while let Some(stream) = acceptor.accept().await {
|
||||
tokio::spawn(async move {
|
||||
if let Err(error) = serve_one(stream).await {
|
||||
tracing::debug!("capability stream failed: {error:#}");
|
||||
}
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
async fn serve_one(mut stream: ByteStream) -> Result<()> {
|
||||
let response = match read_message(&mut stream).await? {
|
||||
CapabilityMessage::Get {
|
||||
version: CAPABILITIES_PROTOCOL_VERSION,
|
||||
} => CapabilityMessage::Manifest {
|
||||
manifest: local_manifest(),
|
||||
},
|
||||
CapabilityMessage::Get { version } => CapabilityMessage::Error {
|
||||
message: format!("unsupported capability protocol {version}"),
|
||||
},
|
||||
_ => CapabilityMessage::Error {
|
||||
message: "expected capability request".to_string(),
|
||||
},
|
||||
};
|
||||
write_message(&mut stream, &response).await?;
|
||||
stream.send.finish()?;
|
||||
let _ = tokio::time::timeout(Duration::from_secs(2), stream.send.stopped()).await;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
pub async fn probe_loop(service: Arc<MusicDhtService>, observed: Arc<ObservedVersions>) {
|
||||
let mut interval = tokio::time::interval(PROBE_INTERVAL);
|
||||
loop {
|
||||
interval.tick().await;
|
||||
let peers = service
|
||||
.connected_peers()
|
||||
.into_iter()
|
||||
.chain(
|
||||
service
|
||||
.known_peers()
|
||||
.into_iter()
|
||||
.map(|contact| contact.peer_id),
|
||||
)
|
||||
.collect::<std::collections::BTreeSet<_>>();
|
||||
for peer in peers {
|
||||
if let Err(error) = probe_peer(&service, peer, &observed).await {
|
||||
tracing::trace!(%peer, "peer capability probe unavailable: {error:#}");
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
async fn probe_peer(
|
||||
service: &MusicDhtService,
|
||||
peer: EndpointId,
|
||||
observed: &ObservedVersions,
|
||||
) -> Result<()> {
|
||||
let mut stream = service.open_stream(peer, CAPABILITIES_ALPN).await?;
|
||||
write_message(
|
||||
&mut stream,
|
||||
&CapabilityMessage::Get {
|
||||
version: CAPABILITIES_PROTOCOL_VERSION,
|
||||
},
|
||||
)
|
||||
.await?;
|
||||
stream.send.finish()?;
|
||||
let response = tokio::time::timeout(Duration::from_secs(5), read_message(&mut stream))
|
||||
.await
|
||||
.context("capability request timed out")??;
|
||||
let CapabilityMessage::Manifest { manifest } = response else {
|
||||
anyhow::bail!("peer did not return a capability manifest");
|
||||
};
|
||||
manifest.validate()?;
|
||||
{
|
||||
let mut versions = lock(&observed.versions);
|
||||
for (id, version) in manifest.protocols {
|
||||
versions
|
||||
.entry(id)
|
||||
.and_modify(|current| *current = (*current).max(version))
|
||||
.or_insert(version);
|
||||
}
|
||||
}
|
||||
lock(&observed.peers).insert(peer.to_string(), manifest.application_version);
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn lock<T>(mutex: &Mutex<T>) -> MutexGuard<'_, T> {
|
||||
mutex
|
||||
.lock()
|
||||
.unwrap_or_else(std::sync::PoisonError::into_inner)
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
#[test]
|
||||
fn local_manifest_lists_every_player_protocol() {
|
||||
let manifest = local_manifest();
|
||||
for id in [
|
||||
"federation_net",
|
||||
"ticket",
|
||||
"rendezvous",
|
||||
"music_dht",
|
||||
"catalog",
|
||||
"audio",
|
||||
"device_sync",
|
||||
"jam",
|
||||
] {
|
||||
assert!(manifest.protocols.contains_key(id), "missing {id}");
|
||||
}
|
||||
manifest.validate().unwrap();
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn snapshot_reports_only_strictly_newer_versions() {
|
||||
let observed = ObservedVersions::default();
|
||||
lock(&observed.versions).insert("music_dht".to_string(), 99);
|
||||
lock(&observed.versions).insert("jam".to_string(), crate::jam::PROTOCOL_VERSION);
|
||||
let snapshot = ProtocolVersions::snapshot(&observed);
|
||||
assert_eq!(snapshot.newer.len(), 1);
|
||||
assert_eq!(snapshot.newer[0].id, "music_dht");
|
||||
}
|
||||
}
|
||||
@@ -13,6 +13,7 @@
|
||||
//! the network too).
|
||||
|
||||
mod audio;
|
||||
mod capabilities;
|
||||
pub mod catalog;
|
||||
|
||||
use std::collections::{HashMap, VecDeque};
|
||||
@@ -36,6 +37,7 @@ use crate::library::NetworkArtistPreview;
|
||||
use crate::library::models::{ArtistRef, TrackItem};
|
||||
|
||||
pub use audio::{AUDIO_ALPN, DownloadProgress, StreamingStart, TrackMetadata};
|
||||
pub use capabilities::ProtocolVersions;
|
||||
pub use catalog::{CATALOG_ALPN, FedAppearsOn, FedArtistCard, FedCardTrack, FedRelease};
|
||||
|
||||
/// How often the published library is re-synchronized with the local index.
|
||||
@@ -386,6 +388,7 @@ pub struct FedStatus {
|
||||
pub last_sync: Option<String>,
|
||||
pub last_error: Option<String>,
|
||||
pub transport: TransportStatsSnapshot,
|
||||
pub protocols: ProtocolVersions,
|
||||
}
|
||||
|
||||
/// Outcome of preparing a federated track for playback.
|
||||
@@ -426,6 +429,7 @@ pub struct Federation {
|
||||
last_sync: std::sync::Mutex<Option<String>>,
|
||||
last_error: std::sync::Mutex<Option<String>>,
|
||||
transport_stats: Arc<TransportStats>,
|
||||
observed_protocols: Arc<capabilities::ObservedVersions>,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone)]
|
||||
@@ -554,6 +558,7 @@ impl Federation {
|
||||
last_sync: std::sync::Mutex::new(None),
|
||||
last_error: std::sync::Mutex::new(initial_error),
|
||||
transport_stats: Arc::new(TransportStats::default()),
|
||||
observed_protocols: Arc::new(capabilities::ObservedVersions::default()),
|
||||
})
|
||||
}
|
||||
|
||||
@@ -656,6 +661,8 @@ impl Federation {
|
||||
.stream_protocol(crate::devices::SYNC_ALPN)
|
||||
// Capability-scoped shared playback control.
|
||||
.stream_protocol(crate::jam::JAM_ALPN)
|
||||
// Informational application/protocol versions.
|
||||
.schema_independent_stream_protocol(capabilities::CAPABILITIES_ALPN)
|
||||
.build()
|
||||
.map_err(|err| anyhow::anyhow!("invalid federation config: {err}"))?;
|
||||
let (service, mut events) = MusicDhtService::start(config)
|
||||
@@ -728,6 +735,14 @@ impl Federation {
|
||||
Arc::clone(&self.jam),
|
||||
Arc::clone(&service),
|
||||
));
|
||||
let capabilities_acceptor = service
|
||||
.stream_acceptor(capabilities::CAPABILITIES_ALPN)
|
||||
.map_err(|err| anyhow::anyhow!("failed to take capabilities acceptor: {err}"))?;
|
||||
let capabilities_serve_task = tokio::spawn(capabilities::serve(capabilities_acceptor));
|
||||
let capabilities_probe_task = tokio::spawn(capabilities::probe_loop(
|
||||
Arc::clone(&service),
|
||||
Arc::clone(&self.observed_protocols),
|
||||
));
|
||||
|
||||
*guard = Some(Running {
|
||||
service,
|
||||
@@ -742,6 +757,8 @@ impl Federation {
|
||||
device_tick_task,
|
||||
jam_serve_task,
|
||||
jam_poll_task,
|
||||
capabilities_serve_task,
|
||||
capabilities_probe_task,
|
||||
],
|
||||
});
|
||||
self.set_error(None);
|
||||
@@ -880,6 +897,7 @@ impl Federation {
|
||||
network: settings.network_id,
|
||||
last_sync: lock(&self.last_sync).clone(),
|
||||
last_error: lock(&self.last_error).clone(),
|
||||
protocols: ProtocolVersions::snapshot(&self.observed_protocols),
|
||||
..FedStatus::default()
|
||||
};
|
||||
if let Some(running) = guard.as_ref() {
|
||||
|
||||
Reference in New Issue
Block a user