audio: implement random SFX output
This commit is contained in:
+10
-11
@@ -85,10 +85,7 @@ impl Display for Role {
|
||||
|
||||
impl Role {
|
||||
fn is_input(&self) -> bool {
|
||||
match self {
|
||||
Role::Mic => true,
|
||||
_ => false
|
||||
}
|
||||
matches!(self, Role::Mic)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -126,7 +123,7 @@ impl AudioSource {
|
||||
|
||||
fn process(&mut self, scope: &ProcessScope) -> Result<Option<f64>, AudioError> {
|
||||
if self.port.connected_count()? > 0 {
|
||||
let buf: Vec<_> = self.port.as_slice(scope).iter().copied().collect();
|
||||
let buf: Vec<_> = self.port.as_slice(scope).to_vec();
|
||||
self.meter.process_interleaved(&buf);
|
||||
self.sample_sink.blocking_send(buf)?;
|
||||
|
||||
@@ -159,7 +156,11 @@ impl AudioSink {
|
||||
}
|
||||
|
||||
fn process(&mut self, scope: &ProcessScope) -> Result<(), AudioError> {
|
||||
let mut next_outbuf = self.sample_src.try_recv()?;
|
||||
let mut next_outbuf = match self.sample_src.try_recv() {
|
||||
Ok(buf) => buf,
|
||||
Err(tokio::sync::mpsc::error::TryRecvError::Empty) => return Ok(()),
|
||||
Err(err) => return Err(err.into())
|
||||
};
|
||||
self.output_buf.append(&mut next_outbuf);
|
||||
|
||||
if self.port.connected_count()? > 0 && !self.output_buf.is_empty() {
|
||||
@@ -167,9 +168,7 @@ impl AudioSink {
|
||||
let mut next_segment: Vec<f32> = self.output_buf.drain(0..(outbuf.len()).min(self.output_buf.len())).collect();
|
||||
let underrun = outbuf.len() - next_segment.len();
|
||||
if underrun > 0 {
|
||||
for _ in 0..underrun {
|
||||
next_segment.push(0.);
|
||||
}
|
||||
next_segment.extend(std::iter::repeat_n(0., underrun));
|
||||
}
|
||||
|
||||
outbuf.copy_from_slice(&next_segment);
|
||||
@@ -228,7 +227,7 @@ impl NotificationHandler for Notify {
|
||||
}).next();
|
||||
|
||||
if let Some((role, Ok(target_port))) = port_match {
|
||||
let cfg_slot = self.config.connections.entry(*role).or_insert_with(|| Default::default());
|
||||
let cfg_slot = self.config.connections.entry(*role).or_default();
|
||||
|
||||
if are_connected {
|
||||
log::info!("{} connected to {}", role, target_port);
|
||||
@@ -278,7 +277,7 @@ pub async fn start_audio_input() -> (AudioInputControl, AudioInStream, AudioOutS
|
||||
} else {
|
||||
(local_port, &peer)
|
||||
};
|
||||
if let Err(err) = client.connect_ports(&src, &dst) {
|
||||
if let Err(err) = client.connect_ports(src, dst) {
|
||||
log::error!("Failed to reconnect {} to {}: {:?}", role, peer_name, err);
|
||||
} else {
|
||||
log::info!("Reconnected {} to {}", role, peer_name);
|
||||
|
||||
+7
-3
@@ -10,7 +10,7 @@ use futures::StreamExt;
|
||||
|
||||
use ratatui::prelude::*;
|
||||
|
||||
use crate::{artifacts::archive::Archive, audio::start_audio_input, scene::{StageDirection, conversation::ConversationEntry}, tts::start_tts, ui::Ui};
|
||||
use crate::{artifacts::archive::Archive, audio::start_audio_input, scene::{StageDirection, conversation::ConversationEntry}, sfx::start_sfx, tts::start_tts, ui::Ui};
|
||||
|
||||
mod scene;
|
||||
mod events;
|
||||
@@ -21,6 +21,7 @@ mod audio;
|
||||
mod artifacts;
|
||||
mod ui;
|
||||
mod widgets;
|
||||
mod sfx;
|
||||
|
||||
// TODO: We should be able to delete entries from the conversation, or at least go back and edit something I said.
|
||||
// TODO: I want a "mark" command or keyboard shortcut, that inserts a marker into the log, so I know where to come back for the next speaking segment.
|
||||
@@ -131,11 +132,14 @@ async fn main() {
|
||||
SaveData::default()
|
||||
};
|
||||
|
||||
let prediction_ctrl = prediction::conversation_task(saved_session, conversation_src).await;
|
||||
let (audio_ctrl, mic_stream, tts_output, _sfx_output) = start_audio_input().await;
|
||||
let (audio_ctrl, mic_stream, tts_output, sfx_output) = start_audio_input().await;
|
||||
let tts_ctrl = start_tts(tts_output).await;
|
||||
let mut sfx_ctrl = start_sfx(sfx_output).await;
|
||||
sfx_ctrl.play_ambient().await.unwrap();
|
||||
let transcription_ctrl = transcription::start_transcription(mic_stream).await;
|
||||
|
||||
let prediction_ctrl = prediction::conversation_task(saved_session, conversation_src, sfx_ctrl).await;
|
||||
|
||||
let mut app = Ui::new(prediction_ctrl, audio_ctrl, transcription_ctrl, tts_ctrl);
|
||||
|
||||
let mut events = EventStream::new();
|
||||
|
||||
@@ -8,7 +8,7 @@ use serde::{Deserialize, Serialize};
|
||||
use serde_json::{Serializer, Value, ser::CompactFormatter};
|
||||
use tokio::sync::{Mutex, mpsc::{self, UnboundedReceiver, UnboundedSender}};
|
||||
|
||||
use crate::{SaveData, artifacts::{Contents, Track, archive::Archive, mixxx::{MixxxDB, MixxxQuery}, tools::DataSource}, prediction::{character::{Character, CharacterControl, CharacterOutput, character_task}, toolbox::{ArchiveToolbox, StageToolbox}}, scene::{Scene, StageDirection, conversation::{ConversationEntry, Speaker}}};
|
||||
use crate::{SaveData, artifacts::{Contents, Track, archive::Archive, mixxx::{MixxxDB, MixxxQuery}, tools::DataSource}, prediction::{character::{Character, CharacterControl, CharacterOutput, character_task}, toolbox::{ArchiveToolbox, StageToolbox}}, scene::{Scene, StageDirection, conversation::{ConversationEntry, Speaker}}, sfx::SfxControl};
|
||||
use tokio_stream::StreamExt;
|
||||
|
||||
pub mod character;
|
||||
@@ -68,7 +68,7 @@ impl SessionControl {
|
||||
}
|
||||
|
||||
pub async fn changed(&mut self) -> SessionUpdate {
|
||||
self.event_src.recv().await.unwrap()
|
||||
self.event_src.recv().await.expect("Session closed")
|
||||
}
|
||||
}
|
||||
|
||||
@@ -84,7 +84,9 @@ struct Conversation {
|
||||
computer_todo: Arc<Mutex<HashMap<String, bool>>>,
|
||||
archive: Arc<Mutex<Archive>>,
|
||||
current_playlist: Vec<Track>,
|
||||
sys_log_messages: UnboundedReceiver<ConversationEntry>
|
||||
sys_log_messages: UnboundedReceiver<ConversationEntry>,
|
||||
|
||||
sfx: SfxControl
|
||||
}
|
||||
|
||||
impl Conversation {
|
||||
@@ -160,6 +162,7 @@ impl Conversation {
|
||||
self.event_sink.send(SessionUpdate::Responses(next_options)).unwrap();
|
||||
},
|
||||
Speaker::ShipComputer => {
|
||||
self.sfx.play_ambient().await.unwrap();
|
||||
let response: ComputerResponse = serde_json::from_value(value).unwrap();
|
||||
self.insert(ConversationEntry::Spoken(Speaker::ShipComputer, response.message)).await;
|
||||
if response.finished.unwrap_or_default() {
|
||||
@@ -279,7 +282,7 @@ impl Conversation {
|
||||
}
|
||||
}
|
||||
|
||||
pub async fn conversation_task(save_data: SaveData, sys_log_messages: tokio::sync::mpsc::UnboundedReceiver<ConversationEntry>) -> SessionControl {
|
||||
pub async fn conversation_task(save_data: SaveData, sys_log_messages: tokio::sync::mpsc::UnboundedReceiver<ConversationEntry>, sfx: SfxControl ) -> SessionControl {
|
||||
let (input_sink, input_src) = tokio::sync::mpsc::unbounded_channel();
|
||||
let (event_sink, event_src) = tokio::sync::mpsc::unbounded_channel();
|
||||
|
||||
@@ -330,7 +333,8 @@ pub async fn conversation_task(save_data: SaveData, sys_log_messages: tokio::syn
|
||||
archive,
|
||||
current_playlist: vec![],
|
||||
sys_log_messages,
|
||||
computer_todo: shared_todo
|
||||
computer_todo: shared_todo,
|
||||
sfx
|
||||
};
|
||||
|
||||
tokio::spawn(async move {
|
||||
|
||||
+81
@@ -0,0 +1,81 @@
|
||||
use rand::seq::IteratorRandom;
|
||||
use symphonia::core::{formats::{TrackType, probe::Hint}, io::MediaSourceStream};
|
||||
|
||||
use crate::audio::AudioOutStream;
|
||||
|
||||
#[derive(Debug)]
|
||||
pub enum SfxRequest {
|
||||
RandomAmbient
|
||||
}
|
||||
|
||||
#[derive(Debug)]
|
||||
pub struct SfxControl {
|
||||
sink: tokio::sync::mpsc::Sender<SfxRequest>
|
||||
}
|
||||
|
||||
impl SfxControl {
|
||||
pub async fn play_ambient(&mut self) -> Result<(), tokio::sync::mpsc::error::SendError<SfxRequest>> {
|
||||
self.sink.send(SfxRequest::RandomAmbient).await
|
||||
}
|
||||
}
|
||||
|
||||
pub async fn start_sfx(audio_sink: AudioOutStream) -> SfxControl {
|
||||
let (event_sink, mut event_src) = tokio::sync::mpsc::channel(32);
|
||||
tokio::spawn(async move {
|
||||
let sfx_dir = std::path::Path::new("./sfx");
|
||||
|
||||
loop {
|
||||
while let Some(event) = event_src.recv().await {
|
||||
match event {
|
||||
SfxRequest::RandomAmbient => {
|
||||
let avail_files = std::fs::read_dir(sfx_dir).unwrap();
|
||||
let chosen_file = avail_files.choose(&mut rand::rng()).unwrap().unwrap();
|
||||
log::debug!("Queuing ambient sound playback with {:?}", chosen_file);
|
||||
let sfx_fd = std::fs::File::open(chosen_file.path()).unwrap();
|
||||
let mss = MediaSourceStream::new(Box::new(sfx_fd), Default::default());
|
||||
let meta_opts = Default::default();
|
||||
let fmt_opts = Default::default();
|
||||
let mut hint = Hint::new();
|
||||
hint.with_extension(".mp3");
|
||||
let mut format = symphonia::default::get_probe()
|
||||
.probe(&hint, mss, fmt_opts, meta_opts)
|
||||
.expect("Unsupported audio format");
|
||||
let track = format.default_track(TrackType::Audio).expect("No audio track");
|
||||
let dec_opts = Default::default();
|
||||
let mut decoder = symphonia::default::get_codecs()
|
||||
.make_audio_decoder(
|
||||
track.codec_params.as_ref().expect("codec params missing").audio().unwrap(),
|
||||
&dec_opts
|
||||
).expect("Unsupported audio codec");
|
||||
log::debug!("Starting stream");
|
||||
loop {
|
||||
let packet = match format.next_packet() {
|
||||
Ok(Some(packet)) => packet,
|
||||
Ok(None) => break,
|
||||
Err(err) => panic!()
|
||||
};
|
||||
|
||||
match decoder.decode_ref(&packet.as_packet_ref()) {
|
||||
Ok(samples) => {
|
||||
let mut channel_bufs: Vec<f32> = vec![];
|
||||
samples.copy_to_vec_interleaved(&mut channel_bufs);
|
||||
let audio_out_buf: Vec<f32> = channel_bufs.windows(samples.byte_len_per_plane()).map(|channels| {
|
||||
let total_volume = channels.iter().cloned().reduce(|a, b| a + b).unwrap_or_default();
|
||||
total_volume / (channels.len() as f32)
|
||||
}).collect();
|
||||
audio_sink.sink.send(audio_out_buf).await.unwrap();
|
||||
},
|
||||
Err(err) => panic!()
|
||||
}
|
||||
}
|
||||
log::debug!("Playback complete");
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
});
|
||||
|
||||
SfxControl {
|
||||
sink: event_sink
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user