I forgot a few
This commit is contained in:
+1
-83
@@ -2,93 +2,11 @@ use librespot::playback::audio_backend::{Sink, SinkAsBytes, SinkError, SinkResul
|
||||
use librespot::playback::convert::Converter;
|
||||
use librespot::playback::decoder::AudioPacket;
|
||||
use log::error;
|
||||
use std::io::{Stdout, Write};
|
||||
use std::io::Write;
|
||||
use tokio::sync::mpsc::UnboundedSender;
|
||||
|
||||
use crate::ipc;
|
||||
use crate::ipc::packet::IpcPacket;
|
||||
use crate::player::stream::Stream;
|
||||
|
||||
pub struct StdoutSink {
|
||||
client: ipc::Client,
|
||||
output: Option<Box<Stdout>>,
|
||||
}
|
||||
|
||||
impl StdoutSink {
|
||||
pub fn new(client: ipc::Client) -> Self {
|
||||
StdoutSink {
|
||||
client,
|
||||
output: None,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl Sink for StdoutSink {
|
||||
fn start(&mut self) -> SinkResult<()> {
|
||||
if let Err(why) = self.client.send(IpcPacket::StartPlayback) {
|
||||
error!("Failed to send start playback packet: {}", why);
|
||||
return Err(SinkError::ConnectionRefused(why.to_string()));
|
||||
}
|
||||
|
||||
self.output.get_or_insert(Box::new(std::io::stdout()));
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn stop(&mut self) -> SinkResult<()> {
|
||||
if let Err(why) = self.client.send(IpcPacket::StopPlayback) {
|
||||
error!("Failed to send stop playback packet: {}", why);
|
||||
return Err(SinkError::ConnectionRefused(why.to_string()));
|
||||
}
|
||||
|
||||
self
|
||||
.output
|
||||
.take()
|
||||
.ok_or_else(|| SinkError::NotConnected("StdoutSink is not connected".to_string()))?
|
||||
.flush()
|
||||
.map_err(|why| SinkError::OnWrite(why.to_string()))?;
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn write(&mut self, packet: AudioPacket, converter: &mut Converter) -> SinkResult<()> {
|
||||
use zerocopy::AsBytes;
|
||||
|
||||
if let AudioPacket::Samples(samples) = packet {
|
||||
let samples_f32: &[f32] = &converter.f64_to_f32(&samples);
|
||||
|
||||
let resampled = samplerate::convert(
|
||||
44100,
|
||||
48000,
|
||||
2,
|
||||
samplerate::ConverterType::Linear,
|
||||
samples_f32,
|
||||
)
|
||||
.expect("to succeed");
|
||||
|
||||
let samples_i16 =
|
||||
&converter.f64_to_s16(&resampled.iter().map(|v| *v as f64).collect::<Vec<f64>>());
|
||||
|
||||
self.write_bytes(samples_i16.as_bytes())?;
|
||||
}
|
||||
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
|
||||
impl SinkAsBytes for StdoutSink {
|
||||
fn write_bytes(&mut self, data: &[u8]) -> SinkResult<()> {
|
||||
self
|
||||
.output
|
||||
.as_deref_mut()
|
||||
.ok_or_else(|| SinkError::NotConnected("StdoutSink is not connected".to_string()))?
|
||||
.write_all(data)
|
||||
.map_err(|why| SinkError::OnWrite(why.to_string()))?;
|
||||
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
|
||||
pub enum SinkEvent {
|
||||
Start,
|
||||
Stop,
|
||||
|
||||
@@ -1,69 +0,0 @@
|
||||
use std::sync::{Arc, Mutex};
|
||||
|
||||
use ipc_channel::ipc::{self, IpcError, IpcOneShotServer, IpcReceiver, IpcSender, TryRecvError};
|
||||
|
||||
use self::packet::IpcPacket;
|
||||
|
||||
pub mod packet;
|
||||
|
||||
pub struct Server {
|
||||
tx: IpcOneShotServer<IpcSender<IpcPacket>>,
|
||||
rx: IpcOneShotServer<IpcReceiver<IpcPacket>>,
|
||||
}
|
||||
|
||||
impl Server {
|
||||
pub fn create() -> Result<(Self, String, String), IpcError> {
|
||||
let (tx, tx_name) = IpcOneShotServer::new().map_err(IpcError::Io)?;
|
||||
let (rx, rx_name) = IpcOneShotServer::new().map_err(IpcError::Io)?;
|
||||
|
||||
Ok((Self { tx, rx }, tx_name, rx_name))
|
||||
}
|
||||
|
||||
pub fn accept(self) -> Result<Client, IpcError> {
|
||||
let (_, tx) = self.tx.accept().map_err(IpcError::Bincode)?;
|
||||
let (_, rx) = self.rx.accept().map_err(IpcError::Bincode)?;
|
||||
|
||||
Ok(Client::new(tx, rx))
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Clone)]
|
||||
pub struct Client {
|
||||
tx: Arc<Mutex<IpcSender<IpcPacket>>>,
|
||||
rx: Arc<Mutex<IpcReceiver<IpcPacket>>>,
|
||||
}
|
||||
|
||||
impl Client {
|
||||
pub fn new(tx: IpcSender<IpcPacket>, rx: IpcReceiver<IpcPacket>) -> Client {
|
||||
Client {
|
||||
tx: Arc::new(Mutex::new(tx)),
|
||||
rx: Arc::new(Mutex::new(rx)),
|
||||
}
|
||||
}
|
||||
|
||||
pub fn connect(tx_name: impl Into<String>, rx_name: impl Into<String>) -> Result<Self, IpcError> {
|
||||
let (tx, remote_rx) = ipc::channel().map_err(IpcError::Io)?;
|
||||
let (remote_tx, rx) = ipc::channel().map_err(IpcError::Io)?;
|
||||
|
||||
let ttx = IpcSender::connect(tx_name.into()).map_err(IpcError::Io)?;
|
||||
let trx = IpcSender::connect(rx_name.into()).map_err(IpcError::Io)?;
|
||||
|
||||
ttx.send(remote_tx).map_err(IpcError::Bincode)?;
|
||||
trx.send(remote_rx).map_err(IpcError::Bincode)?;
|
||||
|
||||
Ok(Client::new(tx, rx))
|
||||
}
|
||||
|
||||
pub fn send(&self, packet: IpcPacket) -> Result<(), IpcError> {
|
||||
self
|
||||
.tx
|
||||
.lock()
|
||||
.expect("to be able to lock")
|
||||
.send(packet)
|
||||
.map_err(IpcError::Bincode)
|
||||
}
|
||||
|
||||
pub fn try_recv(&self) -> Result<IpcPacket, TryRecvError> {
|
||||
self.rx.lock().expect("to be able to lock").try_recv()
|
||||
}
|
||||
}
|
||||
@@ -1,46 +0,0 @@
|
||||
use serde::{Deserialize, Serialize};
|
||||
|
||||
#[derive(Serialize, Deserialize, Debug, Clone)]
|
||||
pub enum IpcPacket {
|
||||
/// Quit the player process
|
||||
Quit,
|
||||
|
||||
/// Connect to Spotify with the given token and device name
|
||||
Connect(String, String),
|
||||
|
||||
/// Disconnect from Spotify (unused)
|
||||
Disconnect,
|
||||
|
||||
/// Unable to connect to Spotify
|
||||
ConnectError(String),
|
||||
|
||||
/// The audio sink has started writing
|
||||
StartPlayback,
|
||||
|
||||
/// The audio sink has stopped writing
|
||||
StopPlayback,
|
||||
|
||||
/// The current Spotify track was changed
|
||||
TrackChange(String),
|
||||
|
||||
/// Spotify playback was started/resumed
|
||||
Playing(String, u32, u32),
|
||||
|
||||
/// Spotify playback was paused
|
||||
Paused(String, u32, u32),
|
||||
|
||||
/// Sent when the user has switched their Spotify device away from Spoticord
|
||||
Stopped,
|
||||
|
||||
/// Request the player to advance to the next track
|
||||
Next,
|
||||
|
||||
/// Request the player to go back to the previous track
|
||||
Previous,
|
||||
|
||||
/// Request the player to pause playback
|
||||
Pause,
|
||||
|
||||
/// Request the player to resume playback
|
||||
Resume,
|
||||
}
|
||||
@@ -13,7 +13,6 @@ mod audio;
|
||||
mod bot;
|
||||
mod consts;
|
||||
mod database;
|
||||
mod ipc;
|
||||
mod librespot_ext;
|
||||
mod player;
|
||||
mod session;
|
||||
|
||||
@@ -5,7 +5,8 @@ use std::{
|
||||
|
||||
use songbird::input::reader::MediaSource;
|
||||
|
||||
const MAX_SIZE: usize = 1 * 1024 * 1024;
|
||||
// TODO: Find optimal value
|
||||
const MAX_SIZE: usize = 1024 * 1024;
|
||||
|
||||
#[derive(Clone)]
|
||||
pub struct Stream {
|
||||
@@ -45,7 +46,7 @@ impl Write for Stream {
|
||||
let mut buffer = mutex.lock().expect("Mutex was poisoned");
|
||||
|
||||
while buffer.len() + buf.len() > MAX_SIZE {
|
||||
buffer = condvar.wait(buffer).unwrap();
|
||||
buffer = condvar.wait(buffer).expect("Mutex was poisoned");
|
||||
}
|
||||
|
||||
buffer.extend_from_slice(buf);
|
||||
|
||||
@@ -764,6 +764,7 @@ impl SpoticordSession {
|
||||
}
|
||||
|
||||
/// Get the channel id
|
||||
#[allow(dead_code)]
|
||||
pub async fn text_channel_id(&self) -> ChannelId {
|
||||
self.0.read().await.text_channel_id
|
||||
}
|
||||
@@ -777,6 +778,7 @@ impl SpoticordSession {
|
||||
self.0.read().await.call.clone()
|
||||
}
|
||||
|
||||
#[allow(dead_code)]
|
||||
pub async fn http(&self) -> Arc<Http> {
|
||||
self.0.read().await.http.clone()
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user