use crate::audio::{AudioDevice, BUFFER_BACKOFF_MICROS, RealtimePlayer, StreamConfig};
use crate::tui::CaptureBuffer;
use crate::{RealtimeChip, VisualSnapshot};
use parking_lot::Mutex;
use std::collections::VecDeque;
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, AtomicU32, Ordering};
pub struct SnapshotDelayBuffer {
snapshots: VecDeque<VisualSnapshot>,
target_delay: usize,
current_delayed: VisualSnapshot,
}
impl SnapshotDelayBuffer {
pub fn new(audio_buffer_samples: usize, batch_size: usize) -> Self {
let target_delay = (audio_buffer_samples / batch_size).max(1);
Self {
snapshots: VecDeque::with_capacity(target_delay + 2),
target_delay,
current_delayed: VisualSnapshot::default(),
}
}
pub fn push(&mut self, snapshot: VisualSnapshot) {
self.snapshots.push_back(snapshot);
if self.snapshots.len() > self.target_delay
&& let Some(delayed) = self.snapshots.pop_front()
{
self.current_delayed = delayed;
}
}
pub fn get_delayed(&self) -> VisualSnapshot {
self.current_delayed
}
pub fn clear(&mut self) {
self.snapshots.clear();
self.current_delayed = VisualSnapshot::default();
}
}
#[derive(Clone, Copy)]
struct ColorFilter {
enabled: bool,
z1_l: f32,
z2_l: f32,
z1_r: f32,
z2_r: f32,
}
impl ColorFilter {
fn new(enabled: bool) -> Self {
Self {
enabled,
z1_l: 0.0,
z2_l: 0.0,
z1_r: 0.0,
z2_r: 0.0,
}
}
fn process_stereo(&mut self, samples: &mut [f32]) {
if !self.enabled {
return;
}
for chunk in samples.chunks_exact_mut(2) {
let filtered_l = (self.z2_l * 0.25) + (self.z1_l * 0.5) + (chunk[0] * 0.25);
self.z2_l = self.z1_l;
self.z1_l = chunk[0];
chunk[0] = filtered_l;
let filtered_r = (self.z2_r * 0.25) + (self.z1_r * 0.5) + (chunk[1] * 0.25);
self.z2_r = self.z1_r;
self.z1_r = chunk[1];
chunk[1] = filtered_r;
}
}
}
const SAMPLE_BATCH_SIZE: usize = 2048;
pub struct StreamingContext {
pub audio_device: AudioDevice,
pub producer_thread: std::thread::JoinHandle<()>,
pub running: Arc<AtomicBool>,
pub player: Arc<Mutex<Box<dyn RealtimeChip>>>,
pub streamer: Arc<RealtimePlayer>,
pub capture: Option<Arc<Mutex<CaptureBuffer>>>,
pub volume: Arc<AtomicU32>,
pub snapshot_delay: Arc<Mutex<SnapshotDelayBuffer>>,
}
impl StreamingContext {
pub fn start(
player: Box<dyn RealtimeChip>,
config: StreamConfig,
color_filter_enabled: bool,
) -> ym2149_ym_replayer::Result<Self> {
Self::start_internal(player, config, color_filter_enabled, None, true)
}
pub fn start_with_capture(
player: Box<dyn RealtimeChip>,
config: StreamConfig,
color_filter_enabled: bool,
capture: Option<Arc<Mutex<CaptureBuffer>>>,
) -> ym2149_ym_replayer::Result<Self> {
Self::start_internal(player, config, color_filter_enabled, capture, true)
}
pub fn start_paused(
player: Box<dyn RealtimeChip>,
config: StreamConfig,
color_filter_enabled: bool,
capture: Option<Arc<Mutex<CaptureBuffer>>>,
) -> ym2149_ym_replayer::Result<Self> {
Self::start_internal(player, config, color_filter_enabled, capture, false)
}
fn start_internal(
player: Box<dyn RealtimeChip>,
config: StreamConfig,
color_filter_enabled: bool,
capture: Option<Arc<Mutex<CaptureBuffer>>>,
auto_start: bool,
) -> ym2149_ym_replayer::Result<Self> {
let streamer = Arc::new(
RealtimePlayer::new(config)
.map_err(|e| format!("Failed to create realtime player: {e}"))?,
);
let audio_device =
AudioDevice::new(config.sample_rate, config.channels, streamer.get_buffer())
.map_err(|e| format!("Failed to create audio device: {e}"))?;
let player = Arc::new(Mutex::new(player));
let running = Arc::new(AtomicBool::new(true));
let volume = Arc::new(AtomicU32::new(100));
let snapshot_delay = Arc::new(Mutex::new(SnapshotDelayBuffer::new(
config.ring_buffer_size,
SAMPLE_BATCH_SIZE,
)));
let running_clone = Arc::clone(&running);
let player_clone = Arc::clone(&player);
let streamer_clone = Arc::clone(&streamer);
let volume_clone = Arc::clone(&volume);
let snapshot_delay_clone = Arc::clone(&snapshot_delay);
let producer_thread = std::thread::spawn(move || {
run_producer_loop(
player_clone,
streamer_clone,
running_clone,
ColorFilter::new(color_filter_enabled),
auto_start,
volume_clone,
snapshot_delay_clone,
);
});
Ok(StreamingContext {
audio_device,
producer_thread,
running,
player,
streamer,
capture,
volume,
snapshot_delay,
})
}
pub fn set_volume(&self, vol: f32) {
let percentage = (vol.clamp(0.0, 1.0) * 100.0) as u32;
self.volume.store(percentage, Ordering::Relaxed);
}
pub fn replace_player(&self, new_player: Box<dyn RealtimeChip>) {
let mut guard = self.player.lock();
guard.stop();
*guard = new_player;
guard.play();
self.snapshot_delay.lock().clear();
}
pub fn get_delayed_snapshot(&self) -> VisualSnapshot {
self.snapshot_delay.lock().get_delayed()
}
pub fn shutdown(self) {
self.running.store(false, Ordering::Relaxed);
if let Err(e) = self.producer_thread.join() {
eprintln!("Warning: Producer thread panicked during shutdown: {e:?}");
}
self.audio_device.finish();
}
}
fn run_producer_loop(
player: Arc<Mutex<Box<dyn RealtimeChip>>>,
streamer: Arc<RealtimePlayer>,
running: Arc<AtomicBool>,
mut color_filter: ColorFilter,
auto_start: bool,
volume: Arc<AtomicU32>,
snapshot_delay: Arc<Mutex<SnapshotDelayBuffer>>,
) {
let mut sample_buffer = [0.0f32; 4096];
if auto_start {
let mut player = player.lock();
player.play();
}
while running.load(Ordering::Relaxed) {
let batch_size = sample_buffer.len();
let snapshot = {
let mut player = player.lock();
if let Some(reason) = player.unsupported_reason() {
eprintln!("{reason}");
running.store(false, Ordering::Relaxed);
break;
}
player.generate_samples_into_stereo(&mut sample_buffer);
player.visual_snapshot()
};
snapshot_delay.lock().push(snapshot);
color_filter.process_stereo(&mut sample_buffer[..batch_size]);
let vol = volume.load(Ordering::Relaxed) as f32 / 100.0;
if vol < 1.0 {
for sample in sample_buffer[..batch_size].iter_mut() {
*sample *= vol;
}
}
let written = streamer.write_blocking(&sample_buffer[..batch_size]);
if written < batch_size {
std::thread::sleep(std::time::Duration::from_micros(BUFFER_BACKOFF_MICROS));
}
}
}