1use super::*;
2#[cfg(target_os = "linux")]
3use crate::hw::alsa::MidiHub;
4#[cfg(target_os = "macos")]
5use crate::hw::coreaudio::{HwDriver, HwOptions, MidiHub};
6#[cfg(target_os = "windows")]
7use crate::hw::options::HwOptions;
8#[cfg(target_os = "freebsd")]
9use crate::hw::oss::MidiHub;
10#[cfg(target_os = "openbsd")]
11use crate::hw::sndio::{HwDriver, HwOptions, MidiHub};
12#[cfg(target_os = "windows")]
13use crate::hw::wasapi::{self, HwDriver, MidiHub};
14#[cfg(target_os = "macos")]
15use crate::workers::coreaudio_worker::HwWorker;
16#[cfg(target_os = "openbsd")]
17use crate::workers::sndio_worker::HwWorker;
18#[cfg(target_os = "windows")]
19use crate::workers::wasapi_worker::HwWorker;
20use crate::{
21 history::{History, UndoEntry},
22 message::{Action, HwMidiEvent, Message, ProcessTask, SessionSlotState},
23 midi::io::MidiEvent,
24 osc::{OscArg, OscServer, build_error_packet, build_osc_packet},
25 state::State,
26 workers::worker::Worker,
27};
28use std::{
29 collections::{HashMap, VecDeque},
30 net::{SocketAddr, UdpSocket},
31 path::{Path, PathBuf},
32 sync::{Arc, atomic::Ordering},
33 time::{Duration, Instant},
34};
35use tokio::sync::mpsc::{Receiver, Sender, channel};
36use tracing::error;
37
38impl Engine {
39 pub fn state(&self) -> Arc<State> {
40 self.state.clone()
41 }
42
43 pub(crate) fn timing_at_sample(&self, sample: usize) -> (f64, u16, u16) {
44 let bpm = self
45 .tempo_points
46 .iter()
47 .filter(|p| p.sample <= sample)
48 .max_by_key(|p| p.sample)
49 .map(|p| p.bpm)
50 .unwrap_or(self.tempo_bpm)
51 .max(1.0);
52 let (num, den) = self
53 .time_signature_points
54 .iter()
55 .filter(|p| p.sample <= sample)
56 .max_by_key(|p| p.sample)
57 .map(|p| (p.numerator.max(1), p.denominator.max(1)))
58 .unwrap_or((self.tsig_num.max(1), self.tsig_denom.max(1)));
59 (bpm, num, den)
60 }
61
62 pub(crate) fn update_global_tempo_from_map(&mut self) {
63 let (bpm, num, den) = self.timing_at_sample(0);
64 self.tempo_bpm = bpm;
65 self.tsig_num = num;
66 self.tsig_denom = den;
67 }
68
69 pub(crate) fn meter_linear_to_db(peak: f32) -> f32 {
70 if peak <= 1.0e-6 {
71 -90.0
72 } else {
73 (20.0 * peak.log10()).clamp(-90.0, 20.0)
74 }
75 }
76
77 pub(crate) const METER_PUBLISH_INTERVAL: Duration = Duration::from_millis(50);
78 pub(crate) const SESSION_RUNTIME_REPORT_INTERVAL: Duration = Duration::from_millis(50);
79 pub(crate) const TRACK_PROCESS_TIMEOUT: Duration = Duration::from_millis(250);
80 #[cfg(any(target_os = "freebsd", target_os = "linux", target_os = "openbsd"))]
81 pub(crate) const HW_OUT_METER_LINEAR_EPSILON: f32 = 0.0025;
82
83 #[cfg(all(unix, not(target_os = "macos")))]
84 pub(crate) fn session_plugins_dir(&self) -> Option<PathBuf> {
85 self.session_dir.as_ref().map(|d| d.join("plugins"))
86 }
87
88 pub(crate) fn session_audio_dir(&self) -> Option<PathBuf> {
89 self.session_dir.as_ref().map(|d| d.join("audio"))
90 }
91
92 pub(crate) fn session_midi_dir(&self) -> Option<PathBuf> {
93 self.session_dir.as_ref().map(|d| d.join("midi"))
94 }
95
96 pub(crate) fn session_peaks_dir(&self) -> Option<PathBuf> {
97 self.session_dir.as_ref().map(|d| d.join("peaks"))
98 }
99
100 pub(crate) fn ensure_session_subdirs(&self) {
101 if let Some(root) = &self.session_dir {
102 let _ = std::fs::create_dir_all(root.join("plugins"));
103 let _ = std::fs::create_dir_all(root.join("audio"));
104 let _ = std::fs::create_dir_all(root.join("midi"));
105 let _ = std::fs::create_dir_all(root.join("peaks"));
106 }
107 }
108
109 pub fn new(rx: Receiver<Message>, tx: Sender<Message>) -> Self {
110 let (meter_snapshot_producer, _) =
111 crate::triple_buffer::triple_buffer(crate::meter::MeterSnapshot::default());
112 let (transport_snapshot_producer, _) =
113 crate::triple_buffer::triple_buffer(crate::meter::TransportSnapshot::default());
114 let (session_runtime_snapshot_producer, _) =
115 crate::triple_buffer::triple_buffer(crate::meter::SessionRuntimeSnapshot::default());
116 Self::new_with_snapshots(
117 rx,
118 tx,
119 meter_snapshot_producer,
120 transport_snapshot_producer,
121 session_runtime_snapshot_producer,
122 )
123 }
124
125 pub fn new_with_snapshots(
126 rx: Receiver<Message>,
127 tx: Sender<Message>,
128 meter_snapshot_producer: crate::triple_buffer::TripleBufferProducer<
129 crate::meter::MeterSnapshot,
130 >,
131 transport_snapshot_producer: crate::triple_buffer::TripleBufferProducer<
132 crate::meter::TransportSnapshot,
133 >,
134 session_runtime_snapshot_producer: crate::triple_buffer::TripleBufferProducer<
135 crate::meter::SessionRuntimeSnapshot,
136 >,
137 ) -> Self {
138 let state = Arc::new(State::default());
139 let initial_state_snapshot = state.lock().snapshot();
140 let state_snapshot = Arc::new(crate::state::StateSlot::from_pointee(
141 initial_state_snapshot.clone(),
142 ));
143 let collector = basedrop::Collector::new();
146 let hw_ports = Arc::new(arc_swap::ArcSwap::from_pointee(
147 crate::plan_builder::HwPorts {
148 buffer_size: 1024,
149 ..Default::default()
150 },
151 ));
152 let initial_plan =
153 { crate::render_plan::RenderPlan::compile(&initial_state_snapshot, &[], &[], 1024) };
154 let plan_slot = Arc::new(crate::render_plan::PlanSlot::from_pointee(
155 basedrop::Owned::new(&collector.handle(), initial_plan),
156 ));
157 let plan_builder = crate::plan_builder::PlanBuilder::spawn(
158 state_snapshot.clone(),
159 hw_ports.clone(),
160 plan_slot.clone(),
161 collector,
162 );
163 let executor = crate::executor::CycleExecutor::new(plan_slot.clone());
164 Self {
165 rx,
166 tx,
167 clients: vec![],
168 state,
169 state_snapshot,
170 workers: vec![],
171 hw_driver: None,
172 hw_driver_info: None,
173 hw_input_ports: Vec::new(),
174 hw_output_ports: Vec::new(),
175 #[cfg(unix)]
176 jack_runtime: None,
177 midi_hub: Some(MidiHub::default()),
178 hw_worker: None,
179 osc_server: None,
180 osc_reply_socket: None,
181 osc_reply_target: None,
182 pending_hw_midi_events: vec![],
183 pending_hw_midi_events_by_device: HashMap::new(),
184 pending_hw_midi_out_events: vec![],
185 pending_hw_midi_out_events_by_device: vec![],
186 active_hw_notes_by_track: HashMap::new(),
187 active_hw_notes_cycle_start: HashMap::new(),
188 midi_hw_in_routes: vec![],
189 midi_hw_out_routes: vec![],
190 midi_hw_thru_routes: vec![],
191 ready_workers: vec![],
192 pending_requests: VecDeque::new(),
193 awaiting_hwfinished: false,
194 handling_hwfinished: false,
195 transport_panic_flush_pending: false,
196 transport_restart_pending: false,
197 notified_loop_wrap_sample: None,
198 transport_sample: 0,
199 hw_input_latency_frames: 0,
200 hw_output_latency_frames: 0,
201 loop_enabled: false,
202 loop_range_samples: None,
203 metronome_enabled: false,
204 tempo_bpm: 120.0,
205 tsig_num: 4,
206 tsig_denom: 4,
207 tempo_points: vec![crate::message::TempoPoint {
208 sample: 0,
209 bpm: 120.0,
210 }],
211 time_signature_points: vec![crate::message::TimeSignaturePoint {
212 sample: 0,
213 numerator: 4,
214 denominator: 4,
215 }],
216 punch_enabled: false,
217 punch_range_samples: None,
218 audio_recordings: std::collections::HashMap::new(),
219 midi_recordings: std::collections::HashMap::new(),
220 completed_audio_recordings: Vec::new(),
221 completed_midi_recordings: Vec::new(),
222 playing: false,
223 transport_running: false,
224 clip_playback_enabled: true,
225 session_clip_playback_enabled: false,
226 session_transport_sample: 0,
227 session_scene_queue: None,
228 session_scene_queue_length_samples: 0,
229 session_current_scene: None,
230 session_current_scene_previous_scene: None,
231 session_current_scene_start_sample: 0,
232 session_current_scene_length_samples: 0,
233 session_completed_clip_passes: Vec::new(),
234 session_reported_clip_passes: std::collections::HashSet::new(),
235 record_enabled: false,
236 step_recording_enabled: false,
237 session_dir: None,
238 hw_out_level_db: 0.0,
239 hw_out_balance: 0.0,
240 hw_out_muted: false,
241 last_hw_out_meter_publish: None,
242 #[cfg(any(target_os = "freebsd", target_os = "linux", target_os = "openbsd"))]
243 last_hw_out_meter_linear: vec![],
244 hw_out_peak_hold_linear: vec![],
245 #[cfg(any(target_os = "freebsd", target_os = "linux", target_os = "openbsd"))]
246 hw_out_meter_publish_phase: false,
247 last_track_meter_publish: None,
248 last_meter_snapshot_publish: None,
249 last_session_report_publish: None,
250 track_meter_linear_by_track: HashMap::new(),
251 meter_snapshot_producer,
252 transport_snapshot_producer,
253 session_runtime_snapshot_producer,
254 executor,
255 plan_builder,
256 plan_slot,
257 hw_ports,
258 pending_node_jobs: VecDeque::new(),
259 latest_hw_out_meter_db: Arc::new(Vec::new()),
260 latest_track_meter_snapshot: Arc::new(Vec::new()),
261 hw_out_loudness_meter: None,
262 latest_hw_out_lufs: None,
263 history: History::default(),
264 history_group: None,
265 history_suspended: false,
266 offline_bounce_jobs: HashMap::new(),
267 pending_bounce_starts: Vec::new(),
268 bounce_worker_tracks: HashMap::new(),
269 pending_midi_learn: None,
270 pending_global_midi_learn: None,
271 pending_session_midi_learn: None,
272 global_midi_learn_play_pause: None,
273 global_midi_learn_stop: None,
274 global_midi_learn_record_toggle: None,
275 session_midi_learn_slots: HashMap::new(),
276 session_midi_learn_scenes: HashMap::new(),
277 session_midi_learn_stop_track: HashMap::new(),
278 session_midi_learn_stop_all: None,
279 midi_cc_gate: HashMap::new(),
280 modulators: Vec::new(),
281 modulator_values: None,
282 }
283 }
284
285 pub(crate) fn publish_state_snapshot(&self) {
286 let snapshot = self.state.lock().snapshot();
287 self.state_snapshot.store(Arc::new(snapshot));
288 }
289
290 pub(crate) fn hw_driver_cycle_samples(&self) -> Option<usize> {
291 self.hw_driver_info.map(|info| info.cycle_samples)
292 }
293
294 #[cfg(unix)]
295 pub(crate) fn jack_cycle_samples(&self) -> Option<usize> {
296 self.jack_runtime.as_ref().map(|j| j.buffer_size)
297 }
298
299 #[cfg(not(unix))]
300 pub(crate) fn jack_cycle_samples(&self) -> Option<usize> {
301 None
302 }
303
304 pub(crate) fn current_cycle_samples(&self) -> usize {
305 self.hw_driver_cycle_samples()
306 .or_else(|| self.jack_cycle_samples())
307 .unwrap_or(0)
308 }
309
310 pub(crate) fn sample_rate(&self) -> f64 {
311 if let Some(info) = self.hw_driver_info {
312 info.sample_rate as f64
313 } else {
314 #[cfg(unix)]
315 {
316 self.jack_runtime
317 .as_ref()
318 .map(|j| j.sample_rate as f64)
319 .unwrap_or(48_000.0)
320 }
321 #[cfg(not(unix))]
322 {
323 48_000.0
324 }
325 }
326 }
327
328 pub(crate) async fn set_hw_playing(&mut self, playing: bool) {
329 if let Some(worker) = &self.hw_worker {
330 let _ = worker.tx.send(Message::HWSetPlaying(playing)).await;
333 } else if let Some(driver) = self.hw_driver.as_mut() {
334 driver.set_playing(playing);
335 }
336 }
337
338 pub(crate) fn active_transport_sample(&self) -> usize {
339 if self.session_clip_playback_enabled && self.playing {
340 self.session_transport_sample
341 } else {
342 self.transport_sample
343 }
344 }
345
346 pub(crate) fn compute_modulator_values(
347 &self,
348 sample: usize,
349 ) -> Arc<std::collections::HashMap<usize, f32>> {
350 let sample_rate = self.sample_rate();
351 let (bpm, tsig_num, tsig_denom) = self.timing_at_sample(sample);
352 let values: std::collections::HashMap<usize, f32> = self
353 .modulators
354 .iter()
355 .filter(|m| m.enabled)
356 .map(|m| {
357 (
358 m.id,
359 m.value_at(sample, sample_rate, bpm, tsig_num, tsig_denom),
360 )
361 })
362 .collect();
363 Arc::new(values)
364 }
365
366 pub(crate) fn apply_modulators(&mut self, sample: usize) -> Vec<Action> {
367 use crate::modulator::ModulatorTarget;
368 let values = self.compute_modulator_values(sample);
369 self.modulator_values = Some(values.clone());
370 let mut echoes = Vec::new();
371 let mut per_track: HashMap<String, (Option<f32>, Option<f32>)> = HashMap::new();
372 let mut clap_params: HashMap<(String, usize, u32), f64> = HashMap::new();
373 let mut vst3_params: HashMap<(String, usize, u32), f32> = HashMap::new();
374 #[cfg(all(unix, not(target_os = "macos")))]
375 let mut lv2_params: HashMap<(String, usize, u32), f32> = HashMap::new();
376 let mut midi_cc_events: HashMap<String, Vec<MidiEvent>> = HashMap::new();
377
378 let map_f32 = |value: f32, min: f32, max: f32| -> f32 {
379 crate::modulator::map_value(value, min, max)
380 };
381 let map_f64 = |value: f32, min: f64, max: f64| -> f64 {
382 crate::modulator::map_value_f64(value, min, max)
383 };
384
385 for m in &self.modulators {
386 if !m.enabled {
387 continue;
388 }
389 let Some(&value) = values.get(&m.id) else {
390 continue;
391 };
392 for target in &m.targets {
393 match target {
394 ModulatorTarget::TrackVolume {
395 track_name,
396 min,
397 max,
398 } => {
399 let clamped = map_f32(value, *min, *max);
400 per_track.entry(track_name.clone()).or_default().0 = Some(clamped);
401 }
402 ModulatorTarget::TrackBalance {
403 track_name,
404 min,
405 max,
406 } => {
407 let clamped = map_f32(value, *min, *max);
408 per_track.entry(track_name.clone()).or_default().1 = Some(clamped);
409 }
410 ModulatorTarget::HwOutVolume { min, max } => {
411 let clamped = map_f32(value, *min, *max);
412 if (self.hw_out_level_db - clamped).abs() > f32::EPSILON {
413 self.hw_out_level_db = clamped;
414 echoes
415 .push(Action::TrackAutomationLevel("hw:out".to_string(), clamped));
416 }
417 }
418 ModulatorTarget::HwOutBalance { min, max } => {
419 let next = map_f32(value, *min, *max).clamp(-1.0, 1.0);
420 if (self.hw_out_balance - next).abs() > f32::EPSILON {
421 self.hw_out_balance = next;
422 echoes.push(Action::TrackAutomationBalance("hw:out".to_string(), next));
423 }
424 }
425 ModulatorTarget::ClapParameter {
426 track_name,
427 instance_id,
428 param_id,
429 min,
430 max,
431 } => {
432 let param_value = map_f64(value, *min, *max);
433 clap_params
434 .insert((track_name.clone(), *instance_id, *param_id), param_value);
435 }
436 ModulatorTarget::Vst3Parameter {
437 track_name,
438 instance_id,
439 param_id,
440 min,
441 max,
442 } => {
443 let param_value = map_f32(value, *min, *max);
444 vst3_params
445 .insert((track_name.clone(), *instance_id, *param_id), param_value);
446 }
447 #[cfg(all(unix, not(target_os = "macos")))]
448 ModulatorTarget::Lv2Parameter {
449 track_name,
450 instance_id,
451 index,
452 min,
453 max,
454 } => {
455 let param_value = map_f32(value, *min, *max);
456 lv2_params.insert((track_name.clone(), *instance_id, *index), param_value);
457 }
458 ModulatorTarget::MidiCc {
459 track_name,
460 channel,
461 cc,
462 } => {
463 let cc_value = (value * 127.0).round() as u8;
464 midi_cc_events
465 .entry(track_name.clone())
466 .or_default()
467 .push(MidiEvent::new(
468 0,
469 vec![0xB0 | (*channel).min(15), (*cc).min(127), cc_value],
470 ));
471 }
472 }
473 }
474 }
475 let state = self.state_snapshot.load_full();
476 for (track_name, (level, balance)) in per_track {
477 if let Some(level) = level
478 && let Some(track) = state.tracks.get(&track_name).cloned()
479 {
480 let t = track.lock();
481 if (t.level() - level).abs() > f32::EPSILON {
482 t.set_level(level);
483 echoes.push(Action::TrackAutomationLevel(track_name.clone(), level));
484 }
485 }
486 if let Some(balance) = balance
487 && let Some(track) = state.tracks.get(&track_name).cloned()
488 {
489 let t = track.lock();
490 let next = balance.clamp(-1.0, 1.0);
491 if (t.balance() - next).abs() > f32::EPSILON {
492 t.set_balance(next);
493 echoes.push(Action::TrackAutomationBalance(track_name.clone(), next));
494 }
495 }
496 }
497
498 for (track_name, events) in midi_cc_events {
499 if let Some(track) = state.tracks.get(&track_name).cloned() {
500 track.lock().rt.pending_modulator_midi_events.extend(events);
501 }
502 }
503
504 for ((track_name, instance_id, param_id), value) in clap_params {
505 if let Some(track) = state.tracks.get(&track_name).cloned()
506 && track
507 .lock()
508 .set_clap_parameter(instance_id, param_id, value)
509 .is_ok()
510 {
511 echoes.push(Action::TrackSetClapParameter {
512 track_name,
513 instance_id,
514 param_id,
515 value,
516 });
517 }
518 }
519 for ((track_name, instance_id, param_id), value) in vst3_params {
520 if let Some(track) = state.tracks.get(&track_name).cloned()
521 && track
522 .lock()
523 .set_vst3_parameter(instance_id, param_id, value)
524 .is_ok()
525 {
526 echoes.push(Action::TrackSetVst3Parameter {
527 track_name,
528 instance_id,
529 param_id,
530 value,
531 });
532 }
533 }
534 #[cfg(all(unix, not(target_os = "macos")))]
535 for ((track_name, instance_id, index), value) in lv2_params {
536 if let Some(track) = state.tracks.get(&track_name).cloned()
537 && track
538 .lock()
539 .set_lv2_control_value(instance_id, index as usize, f64::from(value))
540 .is_ok()
541 {
542 echoes.push(Action::TrackSetLv2ControlValue {
543 track_name,
544 instance_id,
545 index,
546 value,
547 });
548 }
549 }
550
551 echoes
552 }
553
554 pub(crate) fn session_end_sample(&self) -> usize {
555 self.state
556 .lock()
557 .tracks
558 .values()
559 .map(|track| {
560 let track = track.lock();
561 let audio_end = track
562 .audio
563 .clips()
564 .iter()
565 .map(|clip| clip.end)
566 .max()
567 .unwrap_or(0);
568 let midi_end = track
569 .midi
570 .clips()
571 .iter()
572 .map(|clip| clip.end)
573 .max()
574 .unwrap_or(0);
575 audio_end.max(midi_end)
576 })
577 .max()
578 .unwrap_or(0)
579 }
580
581 pub(crate) fn normalize_transport_sample(&self, sample: usize) -> usize {
582 if self.loop_enabled
583 && let Some((loop_start, loop_end)) = self.loop_range_samples
584 && loop_end > loop_start
585 && sample >= loop_end
586 {
587 let loop_len = loop_end - loop_start;
588 return loop_start + (sample - loop_start) % loop_len;
589 }
590 sample
591 }
592
593 pub(crate) fn scheduled_loop_wrap_for_next_cycle(&self) -> Option<(usize, usize, usize)> {
594 if !self.playing || !self.loop_enabled {
595 return None;
596 }
597 let (loop_start, loop_end) = self.loop_range_samples?;
598 if loop_end <= loop_start || self.transport_sample >= loop_end {
599 return None;
600 }
601 let cycle_samples = self.current_cycle_samples();
602 if cycle_samples == 0 {
603 return None;
604 }
605 let next = self.transport_sample.saturating_add(cycle_samples);
606 if next < loop_end {
607 return None;
608 }
609 let after_frames = loop_end.saturating_sub(self.transport_sample);
610 Some((
611 after_frames,
612 loop_start,
613 self.normalize_transport_sample(next),
614 ))
615 }
616
617 pub(crate) fn cycle_segments(&self, frames: usize) -> Vec<(usize, usize, usize)> {
618 if frames == 0 {
619 return vec![];
620 }
621 if !self.loop_enabled {
622 return vec![(
623 self.transport_sample,
624 self.transport_sample.saturating_add(frames),
625 0,
626 )];
627 }
628 let Some((loop_start, loop_end)) = self.loop_range_samples else {
629 return vec![(
630 self.transport_sample,
631 self.transport_sample.saturating_add(frames),
632 0,
633 )];
634 };
635 if loop_end <= loop_start {
636 return vec![(
637 self.transport_sample,
638 self.transport_sample.saturating_add(frames),
639 0,
640 )];
641 }
642 let mut segments = Vec::new();
643 let mut remaining = frames;
644 let mut out_offset = 0usize;
645 let mut current = self.transport_sample;
646 while remaining > 0 {
647 let take = loop_end.saturating_sub(current).min(remaining);
648 if take == 0 {
649 current = loop_start;
650 continue;
651 }
652 segments.push((current, current.saturating_add(take), out_offset));
653 out_offset = out_offset.saturating_add(take);
654 remaining -= take;
655 current = if remaining > 0 {
656 loop_start
657 } else {
658 current.saturating_add(take)
659 };
660 }
661 segments
662 }
663
664 pub(crate) fn recording_segments_for_cycle(&self, frames: usize) -> Vec<(usize, usize, usize)> {
665 let segments = self.cycle_segments(frames);
666 let comp = self.hw_input_latency_frames;
667 let segments: Vec<_> = if comp > 0 {
668 segments
669 .into_iter()
670 .map(|(start, end, offset)| {
671 (start.saturating_sub(comp), end.saturating_sub(comp), offset)
672 })
673 .collect()
674 } else {
675 segments
676 };
677 if !self.punch_enabled {
678 return segments;
679 }
680 let Some((punch_start, punch_end)) = self.punch_range_samples else {
681 return vec![];
682 };
683 if punch_end <= punch_start {
684 return vec![];
685 }
686 let mut clipped = Vec::new();
687 for (segment_start, segment_end, frame_offset) in segments {
688 let start = segment_start.max(punch_start);
689 let end = segment_end.min(punch_end);
690 if end <= start {
691 continue;
692 }
693 let clipped_offset = frame_offset.saturating_add(start.saturating_sub(segment_start));
694 clipped.push((start, end, clipped_offset));
695 }
696 clipped
697 }
698
699 pub async fn init(&mut self) {
700 let max_threads = num_cpus::get();
701 for id in 0..max_threads {
702 let (tx, rx) = channel::<Message>(32);
703 let tx_thread = self.tx.clone();
704 let handler = tokio::spawn(async move {
705 let wrk = Worker::new(id, rx, tx_thread, 8);
706 wrk.await.work().await;
707 });
708 let (node_job_tx, mut node_job_rx) = rtrb::RingBuffer::new(64);
709 let (mut node_result_tx, node_result_rx) = rtrb::RingBuffer::new(64);
710 let node_quit = Arc::new(AtomicBool::new(false));
711 let node_quit_thread = node_quit.clone();
712 let node_thread_handle = std::thread::Builder::new()
713 .name(format!("maolan-node-worker-{id}"))
714 .spawn(move || {
715 crate::enable_flush_denormals_to_zero();
716 if let Err(e) = Worker::try_enable_realtime(8) {
717 tracing::warn!(
718 "Node worker {} realtime priority {} not enabled: {}",
719 id,
720 8,
721 e
722 );
723 }
724 while !node_quit_thread.load(std::sync::atomic::Ordering::Acquire) {
725 match node_job_rx.pop() {
726 Ok(job) => {
727 let mut result = Worker::process_node_job_result(id, job);
728 loop {
729 match node_result_tx.push(result) {
730 Ok(()) => break,
731 Err(rtrb::PushError::Full(returned)) => {
732 if node_quit_thread
733 .load(std::sync::atomic::Ordering::Acquire)
734 {
735 break;
736 }
737 result = returned;
738 std::thread::yield_now();
739 }
740 }
741 }
742 }
743 Err(rtrb::PopError::Empty) => {
744 std::thread::park();
745 }
746 }
747 }
748 })
749 .expect("failed to spawn node worker thread");
750 let node_thread = node_thread_handle.thread().clone();
751 std::mem::forget(node_thread_handle);
752 self.workers.push(WorkerData::with_node_mailbox(
753 tx.clone(),
754 handler,
755 node_job_tx,
756 node_result_rx,
757 node_thread,
758 node_quit,
759 ));
760 }
761 }
762
763 pub(crate) async fn notify_clients(&mut self, action: Result<Action, String>) {
764 self.clients.retain(|client| !client.is_closed());
765 for client in self.clients.iter() {
766 if client
767 .send(Message::Response(action.clone()))
768 .await
769 .is_err()
770 {}
771 }
772 if let Some(reply_to) = self.osc_reply_target {
773 match &action {
774 Err(reason) => {
775 self.send_osc_reply(reply_to, &build_error_packet(reason));
776 }
777 Ok(Action::TrackList(names)) => {
778 let args: Vec<OscArg> = names
779 .iter()
780 .map(|name| OscArg::String(name.clone()))
781 .collect();
782 self.send_osc_reply(
783 reply_to,
784 &build_osc_packet("/response/tracks", &"s".repeat(names.len()), &args),
785 );
786 }
787 Ok(Action::TransportState {
788 sample,
789 tempo_bpm,
790 playing,
791 paused: _,
792 tsig_num,
793 tsig_denom,
794 }) => {
795 self.send_osc_reply(
796 reply_to,
797 &build_osc_packet(
798 "/response/transport",
799 "idffii",
800 &[
801 OscArg::Int(*sample as i32),
802 OscArg::Int(if *playing { 1 } else { 0 }),
803 OscArg::Float(*tempo_bpm as f32),
804 OscArg::Float(0.0), OscArg::Int(*tsig_num as i32),
806 OscArg::Int(*tsig_denom as i32),
807 ],
808 ),
809 );
810 }
811 Ok(Action::MeterSnapshot {
812 hw_out_db,
813 track_meters,
814 ..
815 }) => {
816 let mut args: Vec<OscArg> = Vec::new();
817 args.push(OscArg::Int(hw_out_db.len() as i32));
818 for db in hw_out_db.iter() {
819 args.push(OscArg::Float(*db));
820 }
821 args.push(OscArg::Int(track_meters.len() as i32));
822 for (name, channels) in track_meters.iter() {
823 args.push(OscArg::String(name.clone()));
824 args.push(OscArg::Int(channels.len() as i32));
825 for db in channels.iter() {
826 args.push(OscArg::Float(*db));
827 }
828 }
829 let types = args
830 .iter()
831 .map(|a| match a {
832 OscArg::String(_) => 's',
833 OscArg::Int(_) => 'i',
834 OscArg::Float(_) => 'f',
835 })
836 .collect::<String>();
837 self.send_osc_reply(
838 reply_to,
839 &build_osc_packet("/response/meters", &types, &args),
840 );
841 }
842 Ok(Action::TrackPluginGraph {
843 track_name,
844 plugins,
845 connections: _,
846 connectable_connections: _,
847 }) => {
848 let mut args: Vec<OscArg> = vec![OscArg::String(track_name.clone())];
849 args.push(OscArg::Int(plugins.len() as i32));
850 for plugin in plugins.iter() {
851 args.push(OscArg::Int(plugin.instance_id as i32));
852 args.push(OscArg::String(plugin.format.clone()));
853 args.push(OscArg::String(plugin.uri.clone()));
854 args.push(OscArg::String(plugin.name.clone()));
855 args.push(OscArg::Int(plugin.bypassed as i32));
856 }
857 let types = args
858 .iter()
859 .map(|a| match a {
860 OscArg::String(_) => 's',
861 OscArg::Int(_) => 'i',
862 OscArg::Float(_) => 'f',
863 })
864 .collect::<String>();
865 self.send_osc_reply(
866 reply_to,
867 &build_osc_packet("/response/plugins", &types, &args),
868 );
869 }
870 Ok(Action::ClapPlugins(plugins)) => {
871 let args: Vec<OscArg> = plugins
872 .iter()
873 .map(|p| OscArg::String(format!("{}|{}", p.path, p.name)))
874 .collect();
875 let types = "s".repeat(args.len());
876 self.send_osc_reply(
877 reply_to,
878 &build_osc_packet("/response/clap_plugins", &types, &args),
879 );
880 }
881 Ok(Action::Vst3Plugins(plugins)) => {
882 let args: Vec<OscArg> = plugins
883 .iter()
884 .map(|p| OscArg::String(format!("{}|{}", p.id, p.name)))
885 .collect();
886 let types = "s".repeat(args.len());
887 self.send_osc_reply(
888 reply_to,
889 &build_osc_packet("/response/vst3_plugins", &types, &args),
890 );
891 }
892 #[cfg(all(unix, not(target_os = "macos")))]
893 Ok(Action::Lv2Plugins(plugins)) => {
894 let args: Vec<OscArg> = plugins
895 .iter()
896 .map(|p| OscArg::String(format!("{}|{}", p.uri, p.name)))
897 .collect();
898 let types = "s".repeat(args.len());
899 self.send_osc_reply(
900 reply_to,
901 &build_osc_packet("/response/lv2_plugins", &types, &args),
902 );
903 }
904 Ok(Action::ClapPluginsUnavailable { error })
905 | Ok(Action::Vst3PluginsUnavailable { error }) => {
906 self.send_osc_reply(reply_to, &build_error_packet(error));
907 }
908 #[cfg(all(unix, not(target_os = "macos")))]
909 Ok(Action::Lv2PluginsUnavailable { error }) => {
910 self.send_osc_reply(reply_to, &build_error_packet(error));
911 }
912 Ok(Action::TrackClapParameters {
913 track_name,
914 instance_id,
915 parameters,
916 }) => {
917 let json = serde_json::json!(
918 parameters
919 .iter()
920 .map(|p| serde_json::json!({
921 "id": p.id,
922 "name": p.name,
923 "module": p.module,
924 "min_value": p.min_value,
925 "max_value": p.max_value,
926 "default_value": p.default_value,
927 }))
928 .collect::<Vec<_>>()
929 )
930 .to_string();
931 self.send_osc_reply(
932 reply_to,
933 &build_osc_packet(
934 "/response/plugin_parameters",
935 "siss",
936 &[
937 OscArg::String(track_name.clone()),
938 OscArg::Int(*instance_id as i32),
939 OscArg::String("clap".to_string()),
940 OscArg::String(json),
941 ],
942 ),
943 );
944 }
945 Ok(Action::TrackVst3Parameters {
946 track_name,
947 instance_id,
948 parameters,
949 }) => {
950 let json = serde_json::to_string(parameters).unwrap_or_default();
951 self.send_osc_reply(
952 reply_to,
953 &build_osc_packet(
954 "/response/plugin_parameters",
955 "siss",
956 &[
957 OscArg::String(track_name.clone()),
958 OscArg::Int(*instance_id as i32),
959 OscArg::String("vst3".to_string()),
960 OscArg::String(json),
961 ],
962 ),
963 );
964 }
965 #[cfg(all(unix, not(target_os = "macos")))]
966 Ok(Action::TrackLv2PluginControls {
967 track_name,
968 instance_id,
969 controls,
970 instance_access_handle: _,
971 }) => {
972 let json = serde_json::json!(
973 controls
974 .iter()
975 .map(|c| serde_json::json!({
976 "index": c.index,
977 "name": c.name,
978 "min": c.min,
979 "max": c.max,
980 "value": c.value,
981 }))
982 .collect::<Vec<_>>()
983 )
984 .to_string();
985 self.send_osc_reply(
986 reply_to,
987 &build_osc_packet(
988 "/response/plugin_parameters",
989 "siss",
990 &[
991 OscArg::String(track_name.clone()),
992 OscArg::Int(*instance_id as i32),
993 OscArg::String("lv2".to_string()),
994 OscArg::String(json),
995 ],
996 ),
997 );
998 }
999 Ok(Action::TrackClapNoteNames {
1000 track_name,
1001 note_names,
1002 }) => {
1003 let json = serde_json::to_string(note_names).unwrap_or_default();
1004 self.send_osc_reply(
1005 reply_to,
1006 &build_osc_packet(
1007 "/response/clap_note_names",
1008 "ss",
1009 &[OscArg::String(track_name.clone()), OscArg::String(json)],
1010 ),
1011 );
1012 }
1013 #[cfg(all(unix, not(target_os = "macos")))]
1014 Ok(Action::TrackLv2Midnam {
1015 track_name,
1016 note_names,
1017 }) => {
1018 let json = serde_json::to_string(note_names).unwrap_or_default();
1019 self.send_osc_reply(
1020 reply_to,
1021 &build_osc_packet(
1022 "/response/lv2_midnam",
1023 "ss",
1024 &[OscArg::String(track_name.clone()), OscArg::String(json)],
1025 ),
1026 );
1027 }
1028 _ => {}
1029 }
1030 }
1031 }
1032
1033 pub(crate) fn send_osc_reply(&mut self, reply_to: SocketAddr, packet: &[u8]) {
1034 if self.osc_reply_socket.is_none() {
1035 self.osc_reply_socket = UdpSocket::bind("0.0.0.0:0").ok();
1036 }
1037 if let Some(socket) = self.osc_reply_socket.as_ref() {
1038 let _ = socket.send_to(packet, reply_to);
1039 }
1040 }
1041
1042 pub(crate) async fn dispatch_request(&mut self, a: Action) {
1043 match a {
1044 Action::TrackOfflineBounceCancel { track_name } => {
1045 if let Some(job) = self.offline_bounce_jobs.get(&track_name) {
1046 job.cancel.store(true, Ordering::Relaxed);
1047 }
1048 }
1049 Action::TrackOfflineBounceCancelAll => {
1050 for job in self.offline_bounce_jobs.values() {
1051 job.cancel.store(true, Ordering::Relaxed);
1052 }
1053 }
1054 _ if !self.offline_bounce_jobs.is_empty() => {
1055 self.pending_requests.push_back(a);
1056 }
1057 Action::OpenAudioDevice { .. }
1058 | Action::OpenMidiInputDevice(_)
1059 | Action::OpenMidiOutputDevice(_)
1060 | Action::RequestMeterSnapshot
1061 | Action::RequestTrackList
1062 | Action::RequestTransportState
1063 | Action::Quit
1064 | Action::Log { .. }
1065 | Action::Play
1066 | Action::Pause
1067 | Action::Stop
1068 | Action::TransportPosition(_)
1069 | Action::JumpToEnd
1070 | Action::SetLoopEnabled(_)
1071 | Action::SetLoopRange(_)
1072 | Action::SetPunchEnabled(_)
1073 | Action::SetPunchRange(_)
1074 | Action::SetMetronomeEnabled(_)
1075 | Action::SetTempo(_)
1076 | Action::SetTimeSignature { .. }
1077 | Action::SetTempoMap { .. }
1078 | Action::SetOscEnabled(_)
1079 | Action::SetClipPlaybackEnabled(_)
1080 | Action::SetRecordEnabled(_)
1081 | Action::SetStepRecording(_)
1082 | Action::StepRecordMidiNote { .. }
1083 | Action::SetClipIdentity { .. }
1084 | Action::SetSessionPath(_)
1085 | Action::ClearHistory
1086 | Action::BeginSessionRestore
1087 | Action::PianoKey { .. }
1088 | Action::ModifyMidiNotes { .. }
1089 | Action::ModifyMidiControllers { .. }
1090 | Action::DeleteMidiControllers { .. }
1091 | Action::InsertMidiControllers { .. }
1092 | Action::DeleteMidiNotes { .. }
1093 | Action::InsertMidiNotes { .. }
1094 | Action::SetMidiSysExEvents { .. }
1095 | Action::Session(_) => {
1096 self.handle_request(a).await;
1097 }
1098 #[cfg(all(unix, not(target_os = "macos")))]
1099 Action::ListLv2Plugins => {
1100 self.handle_request(a).await;
1101 }
1102 Action::ListVst3Plugins => {
1103 self.handle_request(a).await;
1104 }
1105 Action::ListClapPlugins => {
1106 self.handle_request(a).await;
1107 }
1108 Action::ListClapPluginsWithCapabilities => {
1109 self.handle_request(a).await;
1110 }
1111 _ => {
1112 self.pending_requests.push_back(a);
1113 if self.can_schedule_hw_cycle() {
1114 self.request_hw_cycle().await;
1115 } else {
1116 while let Some(next) = self.pending_requests.pop_front() {
1117 self.handle_request(next).await;
1118 }
1119 }
1120 }
1121 };
1122 self.publish_clap_state_dirty().await;
1123 }
1124
1125 pub(crate) fn spawn_plugin_host_stderr_reader(
1126 &self,
1127 stderr: std::process::ChildStderr,
1128 source: String,
1129 ) {
1130 let tx = self.tx.clone();
1131 std::thread::spawn(move || {
1132 use std::io::{BufRead, BufReader};
1133 let reader = BufReader::new(stderr);
1134 for line in reader.lines() {
1135 if let Ok(line) = line
1136 && !line.is_empty()
1137 {
1138 let _ = tx.blocking_send(Message::Request(Action::Log {
1139 source: source.clone(),
1140 message: line,
1141 }));
1142 }
1143 }
1144 });
1145 }
1146
1147 pub(crate) fn set_osc_enabled_with<F>(
1148 &mut self,
1149 enabled: bool,
1150 start_server: F,
1151 ) -> Result<(), String>
1152 where
1153 F: FnOnce(Sender<Message>) -> Result<OscServer, String>,
1154 {
1155 if enabled {
1156 if self.osc_server.is_none() {
1157 self.osc_server = Some(start_server(self.tx.clone())?);
1158 }
1159 } else if let Some(mut server) = self.osc_server.take() {
1160 server.stop();
1161 }
1162 Ok(())
1163 }
1164
1165 pub(crate) async fn request_hw_cycle(&mut self) {
1166 if self.awaiting_hwfinished {
1167 tracing::debug!("request_hw_cycle skipped (already awaiting)");
1168 return;
1169 }
1170 tracing::debug!("request_hw_cycle sending TracksFinished");
1171 self.apply_hw_out_gain_and_meter().await;
1172 self.publish_meter_snapshot_if_due();
1173 if let Some((after_frames, loop_start, cycle_end_sample)) =
1174 self.scheduled_loop_wrap_for_next_cycle()
1175 {
1176 self.notified_loop_wrap_sample = Some(cycle_end_sample);
1177 self.notify_clients(Ok(Action::TransportPositionAt {
1178 sample: loop_start,
1179 after_frames,
1180 }))
1181 .await;
1182 } else {
1183 self.notified_loop_wrap_sample = None;
1184 }
1185 if let Some(worker) = &self.hw_worker {
1186 if !self.pending_hw_midi_out_events_by_device.is_empty() {
1187 let out_events = std::mem::take(&mut self.pending_hw_midi_out_events_by_device);
1188 if let Err(e) = worker.tx.send(Message::HWMidiOutEvents(out_events)).await {
1189 error!("Error sending HWMidiOutEvents {e}");
1190 }
1191 }
1192 match worker.tx.send(Message::TracksFinished).await {
1193 Ok(_) => {
1194 self.awaiting_hwfinished = true;
1195 }
1196 Err(e) => {
1197 error!("Error sending TracksFinished {e}");
1198 }
1199 }
1200 }
1201 }
1202
1203 pub(crate) fn should_publish_hw_out_meters(&mut self) -> bool {
1204 let now = Instant::now();
1205 match self.last_hw_out_meter_publish {
1206 Some(last) if now.duration_since(last) < Self::METER_PUBLISH_INTERVAL => false,
1207 _ => {
1208 self.last_hw_out_meter_publish = Some(now);
1209 true
1210 }
1211 }
1212 }
1213
1214 pub(crate) fn should_publish_track_meters(&mut self) -> bool {
1215 let now = Instant::now();
1216 match self.last_track_meter_publish {
1217 Some(last) if now.duration_since(last) < Self::METER_PUBLISH_INTERVAL => false,
1218 _ => {
1219 self.last_track_meter_publish = Some(now);
1220 true
1221 }
1222 }
1223 }
1224
1225 pub(crate) fn should_publish_hw_out_linear(&mut self, peaks_linear: &[f32]) -> bool {
1226 #[cfg(any(target_os = "freebsd", target_os = "linux", target_os = "openbsd"))]
1227 {
1228 self.hw_out_meter_publish_phase = !self.hw_out_meter_publish_phase;
1229 if !self.hw_out_meter_publish_phase {
1230 return false;
1231 }
1232 let changed = if self.last_hw_out_meter_linear.len() != peaks_linear.len() {
1233 true
1234 } else {
1235 self.last_hw_out_meter_linear
1236 .iter()
1237 .zip(peaks_linear.iter())
1238 .any(|(prev, next)| (prev - next).abs() >= Self::HW_OUT_METER_LINEAR_EPSILON)
1239 };
1240 if !changed {
1241 return false;
1242 }
1243 self.last_hw_out_meter_linear.clear();
1244 self.last_hw_out_meter_linear
1245 .extend_from_slice(peaks_linear);
1246 true
1247 }
1248 #[cfg(not(any(target_os = "freebsd", target_os = "linux", target_os = "openbsd")))]
1249 {
1250 let _ = peaks_linear;
1251 false
1252 }
1253 }
1254
1255 pub(crate) async fn maybe_notify_hw_out_meter(&mut self, _meter_db: Vec<f32>) {
1256 {}
1257 }
1258
1259 pub(crate) async fn apply_hw_out_gain_and_meter(&mut self) {
1260 let gain = if self.hw_out_muted {
1261 0.0
1262 } else {
1263 10.0_f32.powf(self.hw_out_level_db / 20.0)
1264 };
1265
1266 if let Some(worker) = &self.hw_worker {
1269 let _ = worker
1270 .tx
1271 .send(Message::HWSetOutputGainBalance {
1272 gain,
1273 balance: self.hw_out_balance,
1274 })
1275 .await;
1276 } else {
1277 #[cfg(unix)]
1278 {
1279 if let Some(jack) = self.jack_runtime.as_ref() {
1280 jack.set_output_gain_linear(gain);
1281 jack.set_output_balance(self.hw_out_balance);
1282 } else {
1283 return;
1284 }
1285 }
1286 #[cfg(not(unix))]
1287 {
1288 return;
1289 }
1290 }
1291
1292 self.feed_hw_out_loudness_meter();
1294
1295 let should_notify_interval = self.should_publish_hw_out_meters();
1296 if !should_notify_interval {
1297 return;
1298 }
1299
1300 let plan = self.executor.plan().clone();
1301 let peaks_linear =
1302 crate::hw::common::output_meter_linear_from_plan(&plan, gain, self.hw_out_balance);
1303 if self.hw_out_peak_hold_linear.len() != peaks_linear.len() {
1304 self.hw_out_peak_hold_linear.resize(peaks_linear.len(), 0.0);
1305 }
1306 let mut held_peaks = Vec::with_capacity(peaks_linear.len());
1307 for (idx, peak_now) in peaks_linear.iter().copied().enumerate() {
1308 let held = self.hw_out_peak_hold_linear[idx] * 0.92;
1309 let next = peak_now.max(held);
1310 self.hw_out_peak_hold_linear[idx] = next;
1311 held_peaks.push(next);
1312 }
1313 let should_notify = self.should_publish_hw_out_linear(&held_peaks);
1314 let meter_db: Vec<f32> = held_peaks
1315 .into_iter()
1316 .map(Self::meter_linear_to_db)
1317 .collect();
1318 self.latest_hw_out_meter_db = Arc::new(meter_db.clone());
1319 if should_notify {
1320 self.maybe_notify_hw_out_meter(meter_db).await;
1321 }
1322 }
1323
1324 fn feed_hw_out_loudness_meter(&mut self) {
1325 let Some(info) = self.hw_driver_info else {
1326 return;
1327 };
1328 let plan = self.executor.plan();
1329 let channels = plan.hw_out_map.len();
1330 if channels == 0 {
1331 return;
1332 }
1333 let sample_rate = info.sample_rate as u32;
1334
1335 let needs_recreate = self
1336 .hw_out_loudness_meter
1337 .as_ref()
1338 .is_none_or(|m| m.channels() != channels || m.sample_rate() != sample_rate);
1339 if needs_recreate {
1340 self.hw_out_loudness_meter =
1341 crate::loudness::LoudnessMeter::new(channels, sample_rate).ok();
1342 }
1343
1344 if let Some(meter) = self.hw_out_loudness_meter.as_mut() {
1345 let interleaved = crate::hw::common::interleaved_hw_out_samples(plan);
1346 meter.feed_interleaved(&interleaved);
1347 }
1348 }
1349
1350 fn update_hw_out_lufs_readout(&mut self) {
1351 self.latest_hw_out_lufs = self
1352 .hw_out_loudness_meter
1353 .as_ref()
1354 .map(|meter| meter.values());
1355 }
1356
1357 pub(crate) fn preload_track_clips_spawn(&self) {
1358 let tracks: Vec<_> = self
1359 .state_snapshot
1360 .load_full()
1361 .tracks
1362 .values()
1363 .cloned()
1364 .collect();
1365 for track in tracks {
1366 tokio::task::spawn_blocking(move || {
1367 track.lock().preload_clips();
1368 });
1369 }
1370 }
1371
1372 pub(crate) fn preload_track_clips(
1381 &self,
1382 ) -> impl std::future::Future<Output = ()> + Send + 'static {
1383 let tracks: Vec<_> = self
1384 .state_snapshot
1385 .load_full()
1386 .tracks
1387 .values()
1388 .cloned()
1389 .collect();
1390 Self::preload_track_handles(tracks)
1391 }
1392
1393 async fn preload_track_handles(tracks: Vec<crate::state::TrackHandle>) {
1394 if tracks.is_empty() {
1395 return;
1396 }
1397 let mut handles = Vec::with_capacity(tracks.len());
1398 for track in tracks {
1399 handles.push(tokio::task::spawn_blocking(move || {
1400 track.lock().preload_clips();
1401 }));
1402 }
1403 for handle in handles {
1404 if let Err(e) = handle.await {
1405 tracing::warn!("Clip preload task panicked: {e}");
1406 }
1407 }
1408 }
1409
1410 pub(crate) fn task_track_name(task: &ProcessTask) -> String {
1411 match task {
1412 ProcessTask::Track(t) | ProcessTask::FolderInput(t) | ProcessTask::FolderOutput(t) => {
1413 t.lock().name.clone()
1414 }
1415 ProcessTask::Plugin { track, .. } => track.lock().name.clone(),
1416 }
1417 }
1418
1419 pub(crate) fn prepare_task_track(&self, task: &ProcessTask) {
1420 let track = match task {
1421 ProcessTask::Track(t) | ProcessTask::FolderInput(t) | ProcessTask::FolderOutput(t) => t,
1422 ProcessTask::Plugin { track, .. } => track,
1423 };
1424 let mut t = track.lock();
1425 let transport_sample = if self.session_clip_playback_enabled && self.playing {
1426 self.session_transport_sample
1427 } else {
1428 self.transport_sample
1429 };
1430 t.set_transport_sample(transport_sample);
1431 t.set_loop_config(self.loop_enabled, self.loop_range_samples);
1432 t.set_transport_timing(self.tempo_bpm, self.tsig_num, self.tsig_denom);
1433 t.set_clip_playback_enabled(self.clip_playback_enabled && self.playing);
1434 t.set_session_clip_playback_enabled(self.session_clip_playback_enabled && self.playing);
1435 t.set_record_tap_enabled(self.playing && self.record_enabled);
1436 t.audio.set_processing(true);
1437 }
1438
1439 pub(crate) async fn dispatch_node_jobs(&mut self, jobs: Vec<crate::executor::NodeJob>) -> bool {
1443 self.pending_node_jobs.extend(jobs);
1444 let mut cycle_complete = false;
1445 while !self.pending_node_jobs.is_empty() {
1446 let Some(worker_index) = self.take_ready_worker_index() else {
1447 break;
1448 };
1449 let Some(job) = self.pending_node_jobs.pop_front() else {
1450 break;
1451 };
1452 if let Some(crate::render_plan::Op::Task { task, .. }) =
1453 job.plan.nodes.get(job.node as usize)
1454 {
1455 self.prepare_task_track(task);
1456 }
1457 let worker = &mut self.workers[worker_index];
1458 if let Some(node_job_tx) = worker.node_job_tx.as_mut() {
1459 match node_job_tx.push(job) {
1460 Ok(()) => {
1461 if let Some(thread) = &worker.node_thread {
1462 thread.unpark();
1463 }
1464 }
1465 Err(rtrb::PushError::Full(job)) => {
1466 self.pending_node_jobs.push_front(job);
1467 self.push_ready_worker(worker_index);
1468 break;
1469 }
1470 }
1471 } else {
1472 let node = job.node;
1473 error!("Worker {worker_index} has no node-job mailbox");
1474 let outcome = self.executor.abandon_node(node, Instant::now());
1475 self.log_silenced_nodes(&outcome.silenced);
1476 cycle_complete |= outcome.cycle_complete;
1477 self.pending_node_jobs.extend(outcome.jobs);
1478 }
1479 }
1480 cycle_complete
1481 }
1482
1483 fn ensure_metronome_wiring(&mut self) {
1491 let Some(track) = self
1492 .state_snapshot
1493 .load_full()
1494 .tracks
1495 .get(Self::METRONOME_TRACK)
1496 .cloned()
1497 else {
1498 return;
1499 };
1500 let frames = self.current_cycle_samples();
1501 let (_, changed) = track.lock().ensure_metronome_source(frames);
1502 if changed {
1503 self.plan_builder.mark_dirty();
1504 }
1505 }
1506
1507 pub(crate) async fn start_plan_cycle(&mut self) -> bool {
1508 if !self.playing || !self.executor.cycle_complete() || !self.offline_bounce_jobs.is_empty()
1512 {
1513 return false;
1514 }
1515 self.refresh_realtime_infection();
1516 self.ensure_metronome_wiring();
1517 let jobs = self.executor.start_cycle(Instant::now());
1518 if self.dispatch_node_jobs(jobs).await {
1519 self.on_all_tracks_finished().await;
1520 return true;
1521 }
1522 false
1523 }
1524
1525 pub(crate) async fn on_node_done(
1529 &mut self,
1530 worker_id: usize,
1531 epoch: u64,
1532 node: u32,
1533 output_linear: Vec<f32>,
1534 parameter_updates: Vec<Action>,
1535 latency_changed: bool,
1536 ) {
1537 self.push_ready_worker(worker_id);
1538 let mut complete = self.dispatch_node_jobs(Vec::new()).await;
1539 if epoch != self.executor.epoch() {
1540 tracing::debug!(
1541 "dropping stale NodeDone (epoch {} vs {}) for node {}",
1542 epoch,
1543 self.executor.epoch(),
1544 node
1545 );
1546 return;
1547 }
1548 if latency_changed {
1549 self.publish_state_snapshot();
1550 self.plan_builder.mark_dirty();
1551 }
1552 let plan = self.executor.plan().clone();
1553 if let Some(crate::render_plan::Op::Task { task, .. }) = plan.nodes.get(node as usize) {
1554 let track_name = Self::task_track_name(task);
1555 self.track_meter_linear_by_track
1556 .insert(track_name, output_linear);
1557 }
1558 for action in parameter_updates {
1559 self.notify_clients(Ok(action)).await;
1560 }
1561 let now = Instant::now();
1562 let (jobs, done) = self.executor.on_node_done(epoch, node, now);
1563 complete |= done;
1564 complete |= self.dispatch_node_jobs(jobs).await;
1565 let outcome = self
1566 .executor
1567 .force_timeouts(now, Self::TRACK_PROCESS_TIMEOUT);
1568 self.log_silenced_nodes(&outcome.silenced);
1569 complete |= self.dispatch_node_jobs(outcome.jobs).await;
1570 complete |= outcome.cycle_complete;
1571 if complete {
1572 self.on_all_tracks_finished().await;
1573 }
1574 }
1575
1576 pub(crate) async fn poll_node_worker_results(&mut self) {
1577 let mut results = Vec::new();
1578 for worker in &mut self.workers {
1579 if let Some(rx) = worker.node_result_rx.as_mut() {
1580 while let Ok(result) = rx.pop() {
1581 results.push(result);
1582 }
1583 }
1584 }
1585 for result in results {
1586 self.on_node_done(
1587 result.worker_id,
1588 result.epoch,
1589 result.node,
1590 result.output_linear,
1591 result.parameter_updates,
1592 result.latency_changed,
1593 )
1594 .await;
1595 }
1596 }
1597
1598 pub(crate) async fn poll_jack_hw_finished(&mut self) {
1599 #[cfg(unix)]
1600 {
1601 let finished = self
1602 .jack_runtime
1603 .as_ref()
1604 .map(|jack| jack.take_hw_finished_count())
1605 .unwrap_or(0);
1606 if finished > 0 {
1607 self.handle_hw_finished().await;
1608 }
1609 }
1610 }
1611
1612 pub(crate) async fn handle_hw_finished(&mut self) {
1613 if !self.awaiting_hwfinished {
1614 tracing::debug!("HWFinished ignored (not awaiting)");
1615 return;
1616 }
1617 tracing::debug!("HWFinished handling; playing={}", self.playing);
1618 self.handling_hwfinished = true;
1619 self.awaiting_hwfinished = false;
1620 #[cfg(unix)]
1621 {
1622 if let Some(jack) = self.jack_runtime.as_mut() {
1623 if !self.pending_hw_midi_out_events.is_empty() {
1624 let out_events = std::mem::take(&mut self.pending_hw_midi_out_events);
1625 jack.write_events(&out_events);
1626 }
1627 let mut in_events = vec![];
1628 jack.read_events_into(&mut in_events);
1629 if !in_events.is_empty() {
1630 self.pending_hw_midi_events.extend(in_events);
1631 }
1632 let dropped = jack.take_midi_events_dropped();
1633 if dropped > 0 {
1634 tracing::warn!(
1635 "JACK MIDI ring full; {dropped} events dropped since last cycle"
1636 );
1637 }
1638 }
1639 }
1640 #[cfg(unix)]
1641 if self.jack_runtime.is_some() {
1642 self.sync_from_jack_transport().await;
1643 }
1644 while let Some(a) = self.pending_requests.pop_front() {
1645 self.handle_request(a).await;
1646 }
1647 self.apply_mute_solo_policy();
1648 self.append_recorded_cycle();
1649 self.flush_completed_recordings().await;
1650 let hw_in_routes = self.midi_hw_in_routes.clone();
1651 let pending_hw_in_by_device = self.pending_hw_midi_events_by_device.clone();
1652 let mut reconfigured_tracks = Vec::new();
1653 let state = self.state_snapshot.load_full();
1654 for (track_name, track) in state.tracks.iter() {
1655 let mut track_lock = track.lock();
1656 if self.jack_runtime_is_some() {
1657 if !self.pending_hw_midi_events.is_empty() {
1658 track_lock.push_hw_midi_events(&self.pending_hw_midi_events);
1659 }
1660 } else {
1661 for route in hw_in_routes.iter().filter(|r| &r.to_track == track_name) {
1662 if let Some(events) = pending_hw_in_by_device.get(&route.device) {
1663 track_lock.push_hw_midi_events_to_port(route.to_port, events);
1664 }
1665 }
1666 }
1667 if track_lock.setup() {
1668 reconfigured_tracks.push(track_name.clone());
1669 }
1670 }
1671 self.publish_track_meters();
1672 self.publish_session_runtime_reports().await;
1673 self.publish_clap_state_dirty().await;
1674 for track_name in reconfigured_tracks {
1675 let track = state.tracks.get(&track_name).cloned();
1676 if let Some(track) = track {
1677 let (plugins, connections, connectable_connections) = {
1678 let track_lock = track.lock();
1679 (
1680 track_lock.plugin_graph_plugins(false),
1681 track_lock.plugin_graph_connections(),
1682 track_lock.connectable_connections(),
1683 )
1684 };
1685 self.notify_clients(Ok(Action::TrackPluginGraph {
1686 track_name: track_name.clone(),
1687 plugins,
1688 connections,
1689 connectable_connections,
1690 }))
1691 .await;
1692 }
1693 }
1694 self.pending_hw_midi_events.clear();
1695 self.pending_hw_midi_events_by_device.clear();
1696 if self.transport_running {
1697 if self.transport_panic_flush_pending {
1698 self.transport_panic_flush_pending = false;
1699 } else if self.transport_restart_pending {
1700 self.transport_restart_pending = false;
1701 } else {
1702 let next = self
1703 .transport_sample
1704 .saturating_add(self.current_cycle_samples());
1705 let normalized = self.normalize_transport_sample(next);
1706 let wrapped = normalized != next;
1707 self.transport_sample = normalized;
1708 self.publish_transport_snapshot();
1709 if wrapped {
1710 if self.notified_loop_wrap_sample == Some(self.transport_sample) {
1711 self.notified_loop_wrap_sample = None;
1712 } else {
1713 self.notify_clients(Ok(Action::TransportPosition(self.transport_sample)))
1714 .await;
1715 }
1716 }
1717 }
1718 }
1719 if self.session_clip_playback_enabled && self.playing {
1720 self.session_transport_sample = self
1721 .session_transport_sample
1722 .saturating_add(self.current_cycle_samples());
1723 }
1724 {
1725 let echoes = self.apply_modulators(self.active_transport_sample());
1726 for action in echoes {
1727 self.notify_clients(Ok(action)).await;
1728 }
1729 }
1730 self.start_plan_cycle().await;
1731 #[cfg(unix)]
1732 {
1733 if self.jack_runtime.is_some() {
1734 self.awaiting_hwfinished = true;
1735 }
1736 }
1737 self.handling_hwfinished = false;
1738 }
1739
1740 pub(crate) async fn on_executor_tick(&mut self) {
1743 if self.executor.cycle_complete() {
1744 return;
1745 }
1746 let outcome = self
1747 .executor
1748 .force_timeouts(Instant::now(), Self::TRACK_PROCESS_TIMEOUT);
1749 self.log_silenced_nodes(&outcome.silenced);
1750 let mut complete = self.dispatch_node_jobs(outcome.jobs).await;
1751 complete |= outcome.cycle_complete;
1752 if complete {
1753 self.on_all_tracks_finished().await;
1754 }
1755 }
1756
1757 pub(crate) fn log_silenced_nodes(&self, silenced: &[u32]) {
1758 for &node in silenced {
1759 let plan = self.executor.plan();
1760 let name = match plan.nodes.get(node as usize) {
1761 Some(crate::render_plan::Op::Task { task, .. }) => Self::task_track_name(task),
1762 _ => format!("node {node}"),
1763 };
1764 tracing::warn!(
1765 "Node {} ('{}') exceeded process timeout ({} ms); forced silent completion for cycle",
1766 node,
1767 name,
1768 Self::TRACK_PROCESS_TIMEOUT.as_millis()
1769 );
1770 }
1771 }
1772
1773 pub(crate) async fn send_bounce_job(
1776 &mut self,
1777 worker_index: usize,
1778 job: crate::message::OfflineBounceWork,
1779 ) {
1780 let track_name = job.track_name.clone();
1781 self.bounce_worker_tracks
1782 .insert(worker_index, track_name.clone());
1783 let worker = &self.workers[worker_index];
1784 if let Err(e) = worker.tx.send(Message::ProcessOfflineBounce(job)).await {
1785 self.bounce_worker_tracks.remove(&worker_index);
1786 self.offline_bounce_jobs.remove(&track_name);
1787 self.push_ready_worker(worker_index);
1788 self.notify_clients(Err(format!("Failed to schedule offline bounce: {e}")))
1789 .await;
1790 }
1791 }
1792
1793 async fn drain_pending_requests_if_idle(&mut self) {
1796 if self.offline_bounce_jobs.is_empty() {
1797 while let Some(next) = self.pending_requests.pop_front() {
1798 self.handle_request(next).await;
1799 }
1800 }
1801 }
1802
1803 pub(crate) async fn on_all_tracks_finished(&mut self) {
1804 let pending = std::mem::take(&mut self.pending_bounce_starts);
1807 for (worker_index, job) in pending {
1808 self.send_bounce_job(worker_index, job).await;
1809 }
1810 if self.transport_restart_pending {
1811 let state = self.state_snapshot.load_full();
1812 for track in state.tracks.values() {
1813 track.lock().take_hw_midi_out_events();
1814 }
1815 } else if self.hw_worker.is_some() {
1816 self.active_hw_notes_cycle_start = self.active_hw_notes_by_track.clone();
1817 let mut out_events = self.collect_hw_midi_output_events_by_device();
1818 if self.loop_enabled
1819 && let Some((_, loop_end)) = self.loop_range_samples
1820 {
1821 let cycle_end = self
1822 .transport_sample
1823 .saturating_add(self.current_cycle_samples());
1824 if self.transport_sample < loop_end && cycle_end >= loop_end {
1825 let wrap_frame = loop_end
1826 .saturating_sub(self.transport_sample)
1827 .min(self.current_cycle_samples())
1828 as u32;
1829 out_events.extend(self.note_off_events_for_active_snapshot(
1830 &self.active_hw_notes_cycle_start,
1831 wrap_frame,
1832 ));
1833 out_events.sort_by(|a, b| {
1834 a.event
1835 .frame
1836 .cmp(&b.event.frame)
1837 .then_with(|| a.device.cmp(&b.device))
1838 });
1839 }
1840 }
1841 self.pending_hw_midi_out_events_by_device.extend(out_events);
1842 } else {
1843 self.pending_hw_midi_out_events = self.collect_hw_midi_output_events();
1844 }
1845 self.request_hw_cycle().await;
1846 }
1847
1848 pub(crate) fn take_ready_worker_index(&mut self) -> Option<usize> {
1849 while !self.ready_workers.is_empty() {
1850 let worker_index = self.ready_workers.remove(0);
1851 if worker_index < self.workers.len() {
1852 return Some(worker_index);
1853 }
1854 }
1855 None
1856 }
1857
1858 pub(crate) fn push_ready_worker(&mut self, worker_index: usize) {
1859 self.ready_workers.push(worker_index);
1860 }
1861
1862 pub(crate) fn publish_track_meters(&mut self) {
1863 if !self.should_publish_track_meters() {
1864 return;
1865 }
1866 let tracks: Vec<(String, crate::state::TrackHandle)> = self
1867 .state_snapshot
1868 .load_full()
1869 .tracks
1870 .iter()
1871 .map(|(name, track)| (name.clone(), track.clone()))
1872 .collect();
1873 let mut snapshot = Vec::with_capacity(tracks.len());
1874 for (name, track) in &tracks {
1875 let linear = self
1876 .track_meter_linear_by_track
1877 .get(name)
1878 .cloned()
1879 .unwrap_or_else(|| track.lock().output_meter_linear());
1880 let output_db = linear
1881 .iter()
1882 .copied()
1883 .map(Self::meter_linear_to_db)
1884 .collect::<Vec<_>>();
1885 snapshot.push((name.clone(), output_db));
1886 }
1887 self.latest_track_meter_snapshot = Arc::new(snapshot);
1888 }
1889
1890 fn record_session_completed_clip_pass(
1891 &mut self,
1892 track_name: String,
1893 scene_index: usize,
1894 clip_id: String,
1895 pass_index: usize,
1896 start_sample: usize,
1897 length_samples: usize,
1898 ) {
1899 let key = (
1900 track_name.clone(),
1901 scene_index,
1902 clip_id.clone(),
1903 pass_index,
1904 start_sample,
1905 );
1906 if self.session_reported_clip_passes.insert(key) {
1907 self.session_completed_clip_passes
1908 .push(crate::meter::SessionCompletedClipPass {
1909 track_name,
1910 scene_index,
1911 clip_id,
1912 pass_index,
1913 start_sample,
1914 length_samples,
1915 });
1916 }
1917 }
1918
1919 fn record_completed_session_scene_span(
1920 &mut self,
1921 tracks: &HashMap<String, crate::state::TrackHandle>,
1922 scene_index: usize,
1923 previous_scene: Option<usize>,
1924 scene_start_sample: usize,
1925 elapsed_samples: usize,
1926 ) {
1927 for (track_name, track) in tracks {
1928 let track = track.lock();
1929 let slot = track.rt.session_slots.get(&scene_index);
1930 let prev_slot = previous_scene.and_then(|scene| track.rt.session_slots.get(&scene));
1931 let playing_clip_id = track
1932 .rt
1933 .playing_session_clips
1934 .last()
1935 .map(|clip| clip.clip_id.as_str());
1936 let clip_id = match (slot, prev_slot) {
1937 (Some(slot), _) if slot.play_enabled => Some(slot.clip_id.as_str()),
1938 (Some(slot), _) if slot.stop_enabled => None,
1939 (_, Some(prev_slot)) if prev_slot.play_enabled => Some(prev_slot.clip_id.as_str()),
1940 (_, Some(prev_slot)) if prev_slot.stop_enabled => None,
1941 _ => playing_clip_id,
1942 }
1943 .filter(|clip_id| !clip_id.is_empty());
1944 let Some(clip_id) = clip_id else { continue };
1945 let clip_length = track
1946 .session_clip_length(clip_id, crate::kind::Kind::Audio)
1947 .or_else(|| track.session_clip_length(clip_id, crate::kind::Kind::MIDI))
1948 .unwrap_or(0);
1949 if clip_length == 0 {
1950 continue;
1951 }
1952 let completed_passes = elapsed_samples / clip_length;
1953 for pass_index in 0..completed_passes {
1954 self.record_session_completed_clip_pass(
1955 track_name.clone(),
1956 scene_index,
1957 clip_id.to_string(),
1958 pass_index,
1959 scene_start_sample.saturating_add(pass_index * clip_length),
1960 clip_length,
1961 );
1962 }
1963 }
1964 }
1965
1966 pub(crate) async fn publish_session_runtime_reports(&mut self) {
1967 let mut current = HashMap::<(String, usize), (SessionSlotState, usize, usize)>::new();
1968 {
1969 let state = self.state_snapshot.load_full();
1970 if let Some((queued_scene, queued_launch_at)) = self.session_scene_queue {
1971 let still_pending = state.tracks.values().any(|track| {
1979 let track = track.lock();
1980 track.rt.pending_session_launches.iter().any(|launch| {
1981 launch.scene_index == queued_scene
1982 && launch.launch_at_sample == queued_launch_at
1983 }) || track
1984 .rt
1985 .playing_session_clips
1986 .iter()
1987 .any(|clip| clip.stop_at_sample == Some(queued_launch_at))
1988 });
1989 if !still_pending && self.session_transport_sample >= queued_launch_at {
1990 if let Some(current_scene) = self.session_current_scene {
1991 let scene_start = self.session_current_scene_start_sample;
1992 let elapsed = queued_launch_at.saturating_sub(scene_start);
1993 self.record_completed_session_scene_span(
1994 &state.tracks,
1995 current_scene,
1996 self.session_current_scene_previous_scene,
1997 scene_start,
1998 elapsed,
1999 );
2000 }
2001 self.session_scene_queue = None;
2002 self.session_current_scene_previous_scene = self.session_current_scene;
2003 self.session_current_scene = Some(queued_scene);
2004 self.session_current_scene_start_sample = queued_launch_at;
2005 self.session_current_scene_length_samples =
2006 self.session_scene_queue_length_samples;
2007 self.session_scene_queue_length_samples = 0;
2008 }
2009 }
2010 if let Some(current_scene) = self.session_current_scene
2011 && self.session_current_scene_length_samples > 0
2012 {
2013 self.record_completed_session_scene_span(
2014 &state.tracks,
2015 current_scene,
2016 self.session_current_scene_previous_scene,
2017 self.session_current_scene_start_sample,
2018 self.session_transport_sample
2019 .saturating_sub(self.session_current_scene_start_sample),
2020 );
2021 }
2022 for (track_name, track) in &state.tracks {
2023 let track = track.lock();
2024 for launch in &track.rt.pending_session_launches {
2025 current.insert(
2026 (track_name.clone(), launch.scene_index),
2027 (SessionSlotState::Queued, 0, 0),
2028 );
2029 }
2030 for clip in &track.rt.playing_session_clips {
2031 if self.session_current_scene_length_samples == 0
2032 && let Some(clip_length) =
2033 track.session_clip_length(&clip.clip_id, clip.kind)
2034 {
2035 let scene_report = self
2036 .session_current_scene
2037 .map(|_| self.session_current_scene_length_samples)
2038 .filter(|length| *length > 0)
2039 .map(|length| {
2040 (
2041 self.session_current_scene.unwrap_or(clip.scene_index),
2042 self.session_current_scene_start_sample,
2043 self.session_transport_sample
2044 .saturating_sub(self.session_current_scene_start_sample),
2045 length,
2046 )
2047 });
2048 let launch_sample = self
2049 .session_transport_sample
2050 .saturating_sub(clip.elapsed_samples);
2051 let (
2052 report_scene_index,
2053 report_start_sample,
2054 report_elapsed,
2055 report_length,
2056 ) = scene_report.unwrap_or((
2057 clip.scene_index,
2058 launch_sample,
2059 clip.elapsed_samples,
2060 clip_length,
2061 ));
2062 if report_length == 0 {
2063 continue;
2064 }
2065 let completed_passes = report_elapsed / report_length;
2066 for pass_index in 0..completed_passes {
2067 self.record_session_completed_clip_pass(
2068 track_name.clone(),
2069 report_scene_index,
2070 clip.clip_id.clone(),
2071 pass_index,
2072 report_start_sample.saturating_add(pass_index * report_length),
2073 report_length,
2074 );
2075 }
2076 }
2077 current.insert(
2078 (track_name.clone(), clip.scene_index),
2079 (
2080 SessionSlotState::Playing,
2081 clip.play_position_samples,
2082 clip.elapsed_samples,
2083 ),
2084 );
2085 }
2086 }
2087 }
2088
2089 if self
2090 .last_session_report_publish
2091 .is_some_and(|t| t.elapsed() < Self::SESSION_RUNTIME_REPORT_INTERVAL)
2092 {
2093 return;
2094 }
2095
2096 let snapshot = self.session_runtime_snapshot_producer.write_buffer();
2097 snapshot.session_sample = self.session_transport_sample;
2098 snapshot.slots.clear();
2099 snapshot.slots.extend(current.iter().map(
2100 |((track_name, scene_index), (state, play_position_samples, elapsed_samples))| {
2101 crate::meter::SessionRuntimeSlotSnapshot {
2102 track_name: track_name.clone(),
2103 scene_index: *scene_index,
2104 state: *state,
2105 play_position_samples: *play_position_samples,
2106 elapsed_samples: *elapsed_samples,
2107 }
2108 },
2109 ));
2110 snapshot.completed_clip_passes = self.session_completed_clip_passes.clone();
2111 snapshot.current_scene = self.session_current_scene;
2112 self.session_runtime_snapshot_producer.publish();
2113 self.last_session_report_publish = Some(Instant::now());
2114 }
2115
2116 pub(crate) async fn publish_clap_state_dirty(&mut self) {
2117 let tracks: Vec<(String, crate::state::TrackHandle)> = self
2118 .state_snapshot
2119 .load_full()
2120 .tracks
2121 .iter()
2122 .map(|(name, track)| (name.clone(), track.clone()))
2123 .collect();
2124 for (track_name, track) in &tracks {
2125 let dirty = track.lock().take_dirty_clap_instances();
2126 for instance_id in dirty {
2127 self.notify_clients(Ok(Action::TrackClapStateDirty {
2128 track_name: track_name.clone(),
2129 instance_id,
2130 }))
2131 .await;
2132 }
2133 }
2134 }
2135
2136 pub(crate) fn reset_meters_after_stop(&mut self) {
2137 self.last_hw_out_meter_publish = None;
2138 self.last_track_meter_publish = None;
2139 self.last_meter_snapshot_publish = None;
2140 self.hw_out_peak_hold_linear.fill(0.0);
2141 #[cfg(any(target_os = "freebsd", target_os = "linux", target_os = "openbsd"))]
2142 {
2143 self.last_hw_out_meter_linear.clear();
2144 }
2145 let hw_channels = self.latest_hw_out_meter_db.len();
2146 self.latest_hw_out_meter_db = Arc::new(vec![-90.0; hw_channels]);
2147
2148 let tracks: Vec<(String, crate::state::TrackHandle)> = self
2149 .state_snapshot
2150 .load_full()
2151 .tracks
2152 .iter()
2153 .map(|(name, track)| (name.clone(), track.clone()))
2154 .collect();
2155 self.track_meter_linear_by_track.clear();
2156 let mut snapshot = Vec::with_capacity(tracks.len());
2157 for (name, track) in tracks {
2158 let mut t = track.lock();
2159 t.clear_output_meters();
2160 let width = t.output_meter_linear().len();
2161 let zero_linear = vec![0.0; width];
2162 self.track_meter_linear_by_track
2163 .insert(name.clone(), zero_linear);
2164 snapshot.push((name, vec![-90.0; width]));
2165 }
2166 self.latest_track_meter_snapshot = Arc::new(snapshot);
2167 self.publish_meter_snapshot();
2168 }
2169
2170 pub(crate) fn publish_meter_snapshot_if_due(&mut self) {
2171 let now = Instant::now();
2172 if self
2173 .last_meter_snapshot_publish
2174 .is_some_and(|last| now.duration_since(last) < Self::METER_PUBLISH_INTERVAL)
2175 {
2176 return;
2177 }
2178 self.last_meter_snapshot_publish = Some(now);
2179 self.publish_meter_snapshot();
2180 }
2181
2182 pub(crate) fn publish_meter_snapshot(&mut self) {
2183 self.update_hw_out_lufs_readout();
2186
2187 let snapshot = self.meter_snapshot_producer.write_buffer();
2188 snapshot.hw_out_db.clear();
2189 snapshot
2190 .hw_out_db
2191 .extend(self.latest_hw_out_meter_db.iter().copied());
2192 snapshot.hw_out_lufs = self.latest_hw_out_lufs;
2193 snapshot.track_meters.clear();
2194 snapshot
2195 .track_meters
2196 .extend(self.latest_track_meter_snapshot.iter().cloned());
2197 self.meter_snapshot_producer.publish();
2198 }
2199
2200 pub(crate) fn publish_transport_snapshot(&mut self) {
2201 let snapshot = self.transport_snapshot_producer.write_buffer();
2202 snapshot.sample = self.transport_sample;
2203 snapshot.tempo_bpm = self.tempo_bpm;
2204 snapshot.playing = self.playing;
2205 snapshot.transport_running = self.transport_running;
2206 snapshot.tsig_num = self.tsig_num;
2207 snapshot.tsig_denom = self.tsig_denom;
2208 self.transport_snapshot_producer.publish();
2209 }
2210
2211 pub(crate) async fn handle_request(&mut self, a: Action) {
2212 match a {
2213 Action::Log { source, message } => {
2214 self.notify_clients(Ok(Action::Log { source, message }))
2215 .await;
2216 }
2217 Action::Undo => {
2218 let actions = match self.history.undo() {
2219 Some(actions) => actions,
2220 None => {
2221 self.notify_clients(Ok(Action::Undo)).await;
2222 self.notify_clients(Ok(Action::HistoryState {
2223 dirty: self.history.is_dirty(),
2224 }))
2225 .await;
2226 return;
2227 }
2228 };
2229
2230 let was_suspended = self.history_suspended;
2231 self.history_suspended = true;
2232 for action in actions {
2233 self.handle_request_inner(action, false).await;
2234 }
2235 self.history_suspended = was_suspended;
2236 self.notify_clients(Ok(Action::Undo)).await;
2237 self.notify_clients(Ok(Action::HistoryState {
2238 dirty: self.history.is_dirty(),
2239 }))
2240 .await;
2241 }
2242 Action::Redo => {
2243 let actions = match self.history.redo() {
2244 Some(actions) => actions,
2245 None => {
2246 self.notify_clients(Ok(Action::Redo)).await;
2247 self.notify_clients(Ok(Action::HistoryState {
2248 dirty: self.history.is_dirty(),
2249 }))
2250 .await;
2251 return;
2252 }
2253 };
2254
2255 let was_suspended = self.history_suspended;
2256 self.history_suspended = true;
2257 for action in actions {
2258 self.handle_request_inner(action, false).await;
2259 }
2260 self.history_suspended = was_suspended;
2261 self.notify_clients(Ok(Action::Redo)).await;
2262 self.notify_clients(Ok(Action::HistoryState {
2263 dirty: self.history.is_dirty(),
2264 }))
2265 .await;
2266 }
2267 Action::ApplyGroupedActions(actions) => {
2268 self.handle_request_inner(Action::BeginHistoryGroup, true)
2269 .await;
2270 for action in actions {
2271 self.handle_request_inner(action, true).await;
2272 }
2273 self.handle_request_inner(Action::EndHistoryGroup, true)
2274 .await;
2275 }
2276 Action::Session(_) => {
2277 self.handle_request_inner(a, false).await;
2278 }
2279 other => {
2280 self.handle_request_inner(other, true).await;
2281 }
2282 }
2283 self.publish_state_snapshot();
2284 }
2285
2286 pub(crate) async fn handle_quit(&mut self, a: Action) {
2287 self.flush_recordings().await;
2288 if let Some(mut worker) = self.hw_worker.take() {
2296 let panic_events = self.panic_events_for_all_hw_midi_outputs();
2299 if !panic_events.is_empty() {
2300 let _ = worker.tx.send(Message::HWMidiOutEvents(panic_events)).await;
2301 }
2302 if let Err(e) = worker.tx.send(Message::Request(a.clone())).await {
2305 error!("Error sending quit message to HW worker: {e}");
2306 }
2307 if let Some(handle) = worker.handle.take() {
2308 handle
2309 .await
2310 .unwrap_or_else(|e| error!("Error waiting for HW worker to quit: {e}"));
2311 }
2312 }
2313 if let Some(hw) = self.hw_driver.as_mut() {
2319 hw.close_fds();
2320 }
2321 if let Some(midi_hub) = self.midi_hub.as_mut() {
2322 midi_hub.close_all();
2323 }
2324 self.hw_driver = None;
2325 self.hw_driver_info = None;
2326 self.hw_input_ports.clear();
2327 self.hw_output_ports.clear();
2328 self.notify_clients(Ok(Action::Quit)).await;
2329 self.ready_workers.clear();
2330 while !self.workers.is_empty() {
2331 let mut worker = self.workers.remove(0);
2332 if let Err(e) = worker.tx.send(Message::Request(a.clone())).await {
2333 error!("Error sending quit message to worker: {e}");
2334 }
2335 if let Some(handle) = worker.handle.take() {
2336 handle
2337 .await
2338 .unwrap_or_else(|e| error!("Error waiting for worker to quit: {e}"));
2339 }
2340 }
2341 #[cfg(unix)]
2342 {
2343 self.jack_runtime = None;
2344 }
2345 self.osc_server = None;
2346 }
2347
2348 #[inline]
2349 pub(crate) fn box_bool<'a>(
2350 fut: impl std::future::Future<Output = bool> + Send + 'a,
2351 ) -> std::pin::Pin<Box<dyn std::future::Future<Output = bool> + Send + 'a>> {
2352 Box::pin(fut)
2353 }
2354
2355 pub(crate) async fn handle_request_inner(
2356 &mut self,
2357 mut action_to_process: Action,
2358 record_history: bool,
2359 ) {
2360 let a = action_to_process.clone();
2361 let suppress_timing_history = self.playing
2362 && matches!(
2363 &action_to_process,
2364 Action::SetTempo(_) | Action::SetTimeSignature { .. } | Action::SetTempoMap { .. }
2365 );
2366 let mut inverse_actions = self.prepare_inverse_actions(
2367 &action_to_process,
2368 record_history,
2369 suppress_timing_history,
2370 );
2371
2372 match action_to_process {
2373 Action::Play => {
2374 if Self::box_bool(self.handle_play(a.clone())).await {
2375 return;
2376 }
2377 }
2378 Action::Pause => {
2379 if Self::box_bool(self.handle_pause(a.clone())).await {
2380 return;
2381 }
2382 }
2383 Action::Stop => {
2384 if Self::box_bool(self.handle_stop(a.clone())).await {
2385 return;
2386 }
2387 }
2388 Action::SessionPlay => {
2389 if Self::box_bool(self.handle_session_play(a.clone())).await {
2390 return;
2391 }
2392 }
2393 Action::JumpToEnd => {
2394 self.transport_sample = self.normalize_transport_sample(self.session_end_sample());
2395 self.publish_transport_snapshot();
2396 self.notify_clients(Ok(Action::TransportPosition(self.transport_sample)))
2397 .await;
2398 }
2399 Action::Panic => {
2400 if Self::box_bool(self.handle_panic(a.clone())).await {
2401 return;
2402 }
2403 }
2404 Action::Session(ref session_action) => {
2405 self.handle_session_action(session_action.clone()).await;
2406 }
2407 Action::SessionRuntimeReport { .. } => {}
2408 Action::SessionMidiLearnTriggered { .. } => {}
2409 Action::SetClipPlaybackEnabled(enabled) => {
2410 self.clip_playback_enabled = enabled;
2411 for track in self.state_snapshot.load_full().tracks.values() {
2412 track.lock().set_clip_playback_enabled(enabled);
2413 }
2414 }
2415 Action::SetSessionClipPlaybackEnabled(enabled) => {
2416 self.session_clip_playback_enabled = enabled;
2417 for track in self.state_snapshot.load_full().tracks.values() {
2418 track.lock().set_session_clip_playback_enabled(enabled);
2419 }
2420 }
2421 Action::TransportPosition(..) => {
2422 if Self::box_bool(self.handle_transport_position(a.clone())).await {
2423 return;
2424 }
2425 }
2426 Action::SetLoopEnabled(enabled) => {
2427 self.loop_enabled = enabled && self.loop_range_samples.is_some();
2428 self.notified_loop_wrap_sample = None;
2429 }
2430 Action::SetLoopRange(..) => {
2431 if Self::box_bool(self.handle_set_loop_range(a.clone())).await {
2432 return;
2433 }
2434 }
2435 Action::SetPunchEnabled(enabled) => {
2436 self.punch_enabled = enabled && self.punch_range_samples.is_some();
2437 }
2438 Action::SetPunchRange(range) => {
2439 self.punch_range_samples = range.and_then(|(start, end)| {
2440 if end > start {
2441 Some((start, end))
2442 } else {
2443 None
2444 }
2445 });
2446 self.punch_enabled = self.punch_range_samples.is_some();
2447 }
2448 Action::SetMetronomeEnabled(enabled) => {
2449 self.metronome_enabled = enabled;
2450 if enabled {
2451 self.ensure_metronome_track().await;
2452 }
2453 if let Some(track) = self
2454 .state_snapshot
2455 .load_full()
2456 .tracks
2457 .get(Self::METRONOME_TRACK)
2458 .cloned()
2459 {
2460 track.lock().set_metronome_enabled(enabled);
2461 }
2462 }
2463 Action::SetTempo(bpm) => {
2464 self.tempo_bpm = bpm.max(1.0);
2465 self.publish_transport_snapshot();
2466 }
2467 Action::SetTimeSignature {
2468 numerator,
2469 denominator,
2470 } => {
2471 self.tsig_num = numerator.max(1);
2472 self.tsig_denom = denominator.max(1);
2473 self.publish_transport_snapshot();
2474 }
2475 Action::SetTempoMap {
2476 ref tempo_points,
2477 ref time_signature_points,
2478 } => {
2479 self.tempo_points = tempo_points.clone();
2480 self.time_signature_points = time_signature_points.clone();
2481 self.update_global_tempo_from_map();
2482 self.publish_transport_snapshot();
2483 }
2484 Action::SetOscEnabled(enabled) => {
2485 if let Err(err) = self.set_osc_enabled_with(enabled, OscServer::start) {
2486 self.notify_clients(Err(err)).await;
2487 }
2488 }
2489 Action::SetRecordEnabled(..) => {
2490 if Self::box_bool(self.handle_set_record_enabled(a.clone())).await {
2491 return;
2492 }
2493 }
2494 Action::SetModulators(ref modulators) => {
2495 self.modulators = modulators.clone();
2496 let echoes = self.apply_modulators(self.active_transport_sample());
2497 for action in echoes {
2498 self.notify_clients(Ok(action)).await;
2499 }
2500 }
2501 Action::SetTrackAutomationLanes {
2502 ref track_name,
2503 ref lanes,
2504 mode,
2505 } => {
2506 if let Some(track) = self.state_snapshot.load_full().tracks.get(track_name) {
2507 let mut track = track.lock();
2508 track.automation_lanes = lanes.clone();
2509 track.set_automation_mode(mode);
2510 }
2511 }
2512 Action::TrackAutomationToggleLane { .. } => {
2513 if Self::box_bool(self.handle_track_automation_toggle_lane(a.clone())).await {
2514 return;
2515 }
2516 }
2517 Action::TrackAutomationInsertPoint { .. } => {
2518 if Self::box_bool(self.handle_track_automation_insert_point(a.clone())).await {
2519 return;
2520 }
2521 }
2522 Action::TrackAutomationDeletePoint { .. } => {
2523 if Self::box_bool(self.handle_track_automation_delete_point(a.clone())).await {
2524 return;
2525 }
2526 }
2527 Action::TrackAutomationSetMode {
2528 ref track_name,
2529 mode,
2530 } => {
2531 if let Some(track) = self
2532 .state_snapshot
2533 .load_full()
2534 .tracks
2535 .get(track_name)
2536 .cloned()
2537 {
2538 track.lock().set_automation_mode(mode);
2539 }
2540 }
2541 Action::RequestTrackList => {
2542 let names: Vec<String> = self
2543 .state_snapshot
2544 .load_full()
2545 .tracks
2546 .keys()
2547 .cloned()
2548 .collect();
2549 self.notify_clients(Ok(Action::TrackList(names))).await;
2550 }
2551 Action::TrackList(_) => {}
2552 Action::RequestTransportState => {
2553 self.notify_clients(Ok(Action::TransportState {
2554 sample: self.transport_sample,
2555 tempo_bpm: self.tempo_bpm,
2556 playing: self.playing,
2557 paused: !self.transport_running && self.playing,
2558 tsig_num: self.tsig_num,
2559 tsig_denom: self.tsig_denom,
2560 }))
2561 .await;
2562 }
2563 Action::TransportState { .. } => {}
2564 Action::SetStepRecording(enabled) => {
2565 self.step_recording_enabled = enabled;
2566 }
2567 Action::BeginHistoryGroup if self.history_group.is_none() => {
2568 self.history_group = Some(UndoEntry {
2569 forward_actions: vec![],
2570 inverse_actions: vec![],
2571 });
2572 }
2573 Action::EndHistoryGroup => {
2574 if let Some(mut group) = self.history_group.take()
2575 && !group.forward_actions.is_empty()
2576 && !group.inverse_actions.is_empty()
2577 {
2578 let mut add_tracks = Vec::new();
2579 let mut connections = Vec::new();
2580 let mut rest = Vec::new();
2581 for action in group.inverse_actions {
2582 if matches!(action, Action::AddTrack { .. }) {
2583 add_tracks.push(action);
2584 } else if matches!(action, Action::Connect { .. }) {
2585 connections.push(action);
2586 } else {
2587 rest.push(action);
2588 }
2589 }
2590 group.inverse_actions = add_tracks;
2591 group.inverse_actions.extend(rest);
2592 group.inverse_actions.extend(connections);
2593 self.history.record(group);
2594 }
2595 }
2596 Action::SetSessionPath(ref path) => {
2597 self.session_dir = Some(Path::new(path).to_path_buf());
2598 self.ensure_session_subdirs();
2599 #[cfg(all(unix, not(target_os = "macos")))]
2600 let _lv2_dir = self.session_plugins_dir();
2601 for track in self.state_snapshot.load_full().tracks.values() {
2602 track.lock().set_session_base_dir(self.session_dir.clone());
2603 }
2604 }
2605 Action::MarkHistorySavePoint => {
2606 self.history.mark_save_point();
2607 self.notify_clients(Ok(Action::HistoryState {
2608 dirty: self.history.is_dirty(),
2609 }))
2610 .await;
2611 }
2612 Action::ClearHistory => {
2613 self.history.clear();
2614 self.history.mark_save_point();
2615 }
2616 Action::BeginSessionRestore => {
2617 self.history_suspended = true;
2618 self.history.clear();
2619 }
2620 Action::EndSessionRestore => {
2621 self.history.clear();
2622 self.history_suspended = false;
2623 self.preload_track_clips_spawn();
2624 }
2625 Action::Quit => {
2626 self.handle_quit(a.clone()).await;
2627 return;
2628 }
2629 Action::AddTrack {
2630 ref name,
2631 audio_ins,
2632 midi_ins,
2633 audio_outs,
2634 midi_outs,
2635 folder,
2636 } => {
2637 self.handle_add_track(
2638 name.clone(),
2639 audio_ins,
2640 midi_ins,
2641 audio_outs,
2642 midi_outs,
2643 folder,
2644 )
2645 .await;
2646 }
2647 Action::TrackAddAudioInput(..) => {
2648 if Self::box_bool(self.handle_track_add_audio_input(a.clone())).await {
2649 return;
2650 }
2651 }
2652 Action::TrackAddAudioOutput(..) => {
2653 if Self::box_bool(self.handle_track_add_audio_output(a.clone())).await {
2654 return;
2655 }
2656 }
2657 Action::TrackRemoveAudioInput(..) => {
2658 if Self::box_bool(self.handle_track_remove_audio_input(a.clone())).await {
2659 return;
2660 }
2661 }
2662 Action::TrackRemoveAudioOutput(..) => {
2663 if Self::box_bool(self.handle_track_remove_audio_output(a.clone())).await {
2664 return;
2665 }
2666 }
2667 Action::RenameTrack { .. } => {
2668 if Self::box_bool(self.handle_rename_track(a.clone())).await {
2669 return;
2670 }
2671 }
2672 Action::RemoveTrack(ref name) => {
2673 self.handle_remove_track(name.clone(), record_history).await;
2674 inverse_actions = None;
2675 }
2676 Action::TrackLevel(ref name, level) => {
2677 if name == "hw:out" {
2678 self.hw_out_level_db = level;
2679 } else if let Some(track) = self.state_snapshot.load_full().tracks.get(name) {
2680 track.lock().set_level(level);
2681 }
2682 }
2683 Action::TrackBalance(ref name, balance) => {
2684 if name == "hw:out" {
2685 self.hw_out_balance = balance.clamp(-1.0, 1.0);
2686 } else if let Some(track) = self.state_snapshot.load_full().tracks.get(name) {
2687 track.lock().set_balance(balance);
2688 }
2689 }
2690 Action::TrackAutomationLevel(ref name, level) => {
2691 tracing::debug!(%name, level, "engine received TrackAutomationLevel");
2692 if name == "hw:out" {
2693 self.hw_out_level_db = level;
2694 } else if let Some(track) = self.state_snapshot.load_full().tracks.get(name) {
2695 track.lock().set_level(level);
2696 }
2697 }
2698 Action::TrackAutomationBalance(ref name, balance) => {
2699 if name == "hw:out" {
2700 self.hw_out_balance = balance.clamp(-1.0, 1.0);
2701 } else if let Some(track) = self.state_snapshot.load_full().tracks.get(name) {
2702 track.lock().set_balance(balance);
2703 }
2704 }
2705 Action::TrackMidiCc { .. } => {
2706 if Self::box_bool(self.handle_track_midi_cc(a.clone())).await {
2707 return;
2708 }
2709 }
2710 Action::RequestMeterSnapshot => {
2711 self.notify_clients(Ok(Action::MeterSnapshot {
2712 hw_out_db: self.latest_hw_out_meter_db.clone(),
2713 hw_out_lufs: self.latest_hw_out_lufs,
2714 track_meters: self.latest_track_meter_snapshot.clone(),
2715 }))
2716 .await;
2717 return;
2718 }
2719 Action::TrackMeters { .. } => {}
2720 Action::MeterSnapshot { .. } => {}
2721 Action::TrackToggleArm(..) => {
2722 if Self::box_bool(self.handle_track_toggle_arm(a.clone())).await {
2723 return;
2724 }
2725 }
2726 Action::TrackToggleMute(ref name) => {
2727 if name == "hw:out" {
2728 self.hw_out_muted = !self.hw_out_muted;
2729 } else if let Some(track) = self.state_snapshot.load_full().tracks.get(name) {
2730 track.lock().mute();
2731 }
2732 }
2733 Action::TrackTogglePhase(ref name) => {
2734 if let Some(track) = self.state_snapshot.load_full().tracks.get(name) {
2735 track.lock().invert_phase();
2736 }
2737 }
2738 Action::TrackToggleSolo(ref name) => {
2739 if name == "hw:out" {
2740 return;
2741 }
2742 if let Some(track) = self.state_snapshot.load_full().tracks.get(name) {
2743 track.lock().solo();
2744 }
2745 }
2746 Action::TrackToggleMaster(ref name) => {
2747 if let Some(track) = self.state_snapshot.load_full().tracks.get(name) {
2748 track.lock().toggle_master();
2749 }
2750 }
2751 Action::TrackToggleInputMonitor {
2752 ref track_name,
2753 lane,
2754 } => {
2755 if let Some(track) = self.state_snapshot.load_full().tracks.get(track_name) {
2756 track.lock().toggle_input_monitor(lane);
2757 }
2758 }
2759 Action::TrackToggleDiskMonitor {
2760 ref track_name,
2761 lane,
2762 } => {
2763 if let Some(track) = self.state_snapshot.load_full().tracks.get(track_name) {
2764 track.lock().toggle_disk_monitor(lane);
2765 }
2766 }
2767 Action::TrackToggleMidiInputMonitor {
2768 ref track_name,
2769 lane,
2770 } => {
2771 if let Some(track) = self.state_snapshot.load_full().tracks.get(track_name) {
2772 track.lock().toggle_midi_input_monitor(lane);
2773 }
2774 }
2775 Action::TrackToggleMidiDiskMonitor {
2776 ref track_name,
2777 lane,
2778 } => {
2779 if let Some(track) = self.state_snapshot.load_full().tracks.get(track_name) {
2780 track.lock().toggle_midi_disk_monitor(lane);
2781 }
2782 }
2783 Action::TrackSetColor {
2784 ref track_name,
2785 color,
2786 } => {
2787 if let Some(track) = self.state_snapshot.load_full().tracks.get(track_name) {
2788 track.lock().color = color;
2789 }
2790 }
2791 Action::TrackArmMidiLearn {
2792 ref track_name,
2793 target,
2794 } => {
2795 if let Err(e) = self.track_handle_or_err(track_name) {
2796 self.notify_clients(Err(e)).await;
2797 return;
2798 }
2799 self.pending_midi_learn = Some((track_name.clone(), target, None));
2800 }
2801 Action::GlobalArmMidiLearn { target } => {
2802 self.pending_global_midi_learn = Some(target);
2803 }
2804 Action::SessionArmMidiLearn { ref target } => {
2805 self.pending_session_midi_learn = Some(target.clone());
2806 }
2807 Action::TrackSetMidiLearnBinding { .. } => {
2808 if Self::box_bool(self.handle_track_set_midi_learn_binding(a.clone())).await {
2809 return;
2810 }
2811 }
2812 Action::SetGlobalMidiLearnBinding { .. } => {
2813 if Self::box_bool(self.handle_set_global_midi_learn_binding(a.clone())).await {
2814 return;
2815 }
2816 }
2817 Action::SetSessionMidiLearnBinding { .. } => {
2818 if Self::box_bool(self.handle_set_session_midi_learn_binding(a.clone())).await {
2819 return;
2820 }
2821 }
2822 Action::TrackSetFolder { .. } => {
2823 if Self::box_bool(self.handle_track_set_folder(a.clone())).await {
2824 return;
2825 }
2826 }
2827 Action::TrackSetParent {
2828 ref track_name,
2829 ref parent_name,
2830 } => {
2831 self.handle_track_set_parent(track_name.as_str(), parent_name.as_deref())
2832 .await;
2833 }
2834 Action::TrackToggleFolder { .. } => {
2835 if Self::box_bool(self.handle_track_toggle_folder(a.clone())).await {
2836 return;
2837 }
2838 }
2839 Action::TrackSetMidiLaneChannel { .. } => {
2840 if Self::box_bool(self.handle_track_set_midi_lane_channel(a.clone())).await {
2841 return;
2842 }
2843 }
2844 Action::TrackSetFrozen { .. } => {
2845 if Self::box_bool(self.handle_track_set_frozen(a.clone())).await {
2846 return;
2847 }
2848 }
2849 Action::TrackSetSessionSlot { .. } => {
2850 if Self::box_bool(self.handle_track_set_session_slot(a.clone())).await {
2851 return;
2852 }
2853 }
2854 Action::TrackSetSessionSlotPlayEnabled { .. } => {
2855 if self
2856 .handle_track_set_session_slot_play_enabled(a.clone())
2857 .await
2858 {
2859 return;
2860 }
2861 }
2862 Action::TrackSetSessionSlotStopEnabled { .. } => {
2863 if self
2864 .handle_track_set_session_slot_stop_enabled(a.clone())
2865 .await
2866 {
2867 return;
2868 }
2869 }
2870 Action::TrackOfflineBounce { .. } => {
2871 self.handle_track_offline_bounce(action_to_process).await;
2872 return;
2873 }
2874 Action::TrackOfflineBounceCancel { .. } => {}
2875 Action::TrackOfflineBounceCancelAll => {}
2876 Action::TrackOfflineBounceCanceled { .. } => {}
2877 Action::TrackOfflineBounceProgress { .. } => {}
2878 Action::PianoKey {
2879 ref track_name,
2880 note,
2881 velocity,
2882 on,
2883 } => {
2884 if let Some(track) = self.state_snapshot.load_full().tracks.get(track_name) {
2885 let status = if on { 0x90 } else { 0x80 };
2886 let event = MidiEvent::new(0, vec![status, note.min(127), velocity.min(127)]);
2887 track.lock().push_hw_midi_events(&[event]);
2888 }
2889 }
2890 Action::ModifyMidiNotes { .. }
2891 | Action::ModifyMidiControllers { .. }
2892 | Action::DeleteMidiControllers { .. }
2893 | Action::InsertMidiControllers { .. }
2894 | Action::DeleteMidiNotes { .. }
2895 | Action::InsertMidiNotes { .. } => {
2896 if let Err(e) = self.apply_midi_edit_action(&action_to_process) {
2897 self.notify_clients(Err(e)).await;
2898 return;
2899 }
2900 }
2901 Action::SetMidiSysExEvents { .. } => {
2902 if let Err(e) = self.apply_midi_edit_action(&action_to_process) {
2903 self.notify_clients(Err(e)).await;
2904 return;
2905 }
2906 }
2907 Action::TrackClearDefaultPassthrough { .. } => {
2908 if Self::box_bool(self.handle_track_clear_default_passthrough(a.clone())).await {
2909 return;
2910 }
2911 }
2912 Action::TrackClearPlugins { .. } => {
2913 if Self::box_bool(self.handle_track_clear_plugins(a.clone())).await {
2914 return;
2915 }
2916 }
2917 #[cfg(all(unix, not(target_os = "macos")))]
2918 Action::ClipSetLv2PluginState { ref track_name, .. } => {
2919 self.notify_clients(Err(format!(
2920 "Track '{}': clip LV2 plugin state changes are not supported",
2921 track_name
2922 )))
2923 .await;
2924 }
2925 Action::TrackGetClapNoteNames { .. } => {
2926 if Self::box_bool(self.handle_track_get_clap_note_names(a.clone())).await {
2927 return;
2928 }
2929 }
2930 #[cfg(all(unix, not(target_os = "macos")))]
2931 Action::TrackGetLv2Midnam { .. } => {
2932 if Self::box_bool(self.handle_track_get_lv2_midnam(a.clone())).await {
2933 return;
2934 }
2935 }
2936 Action::TrackGetPluginGraph { .. } => {
2937 if Self::box_bool(self.handle_track_get_plugin_graph(a.clone())).await {
2938 return;
2939 }
2940 }
2941 Action::TrackPluginGraph { .. } => {}
2942 Action::TrackConnectPluginAudio { .. } => {
2943 if Self::box_bool(self.handle_track_connect_plugin_audio(a.clone())).await {
2944 return;
2945 }
2946 }
2947 Action::TrackConnectPluginMidi { .. } => {
2948 if Self::box_bool(self.handle_track_connect_plugin_midi(a.clone())).await {
2949 return;
2950 }
2951 }
2952 Action::TrackDisconnectPluginAudio { .. } => {
2953 if Self::box_bool(self.handle_track_disconnect_plugin_audio(a.clone())).await {
2954 return;
2955 }
2956 }
2957 Action::TrackDisconnectPluginMidi { .. } => {
2958 if Self::box_bool(self.handle_track_disconnect_plugin_midi(a.clone())).await {
2959 return;
2960 }
2961 }
2962 Action::TrackConnectAudio { .. } => {
2963 if Self::box_bool(self.handle_track_connect_audio(a.clone())).await {
2964 return;
2965 }
2966 }
2967 Action::TrackDisconnectAudio { .. } => {
2968 if Self::box_bool(self.handle_track_disconnect_audio(a.clone())).await {
2969 return;
2970 }
2971 }
2972 Action::TrackConnectMidi { .. } => {
2973 if Self::box_bool(self.handle_track_connect_midi(a.clone())).await {
2974 return;
2975 }
2976 }
2977 Action::TrackDisconnectMidi { .. } => {
2978 if Self::box_bool(self.handle_track_disconnect_midi(a.clone())).await {
2979 return;
2980 }
2981 }
2982 #[cfg(all(unix, not(target_os = "macos")))]
2983 Action::ListLv2Plugins => {
2984 if Self::box_bool(self.handle_list_lv2_plugins(a.clone())).await {
2985 return;
2986 }
2987 }
2988 #[cfg(all(unix, not(target_os = "macos")))]
2989 Action::Lv2Plugins(_) => {}
2990 #[cfg(all(unix, not(target_os = "macos")))]
2991 Action::Lv2PluginsUnavailable { .. } => {}
2992 Action::ListVst3Plugins => {
2993 if Self::box_bool(self.handle_list_vst3_plugins(a.clone())).await {
2994 return;
2995 }
2996 }
2997 Action::Vst3Plugins(_) => {}
2998 Action::Vst3PluginsUnavailable { .. } => {}
2999 Action::ListClapPlugins => {
3000 if Self::box_bool(self.handle_list_clap_plugins(a.clone())).await {
3001 return;
3002 }
3003 }
3004 Action::ListClapPluginsWithCapabilities => {
3005 if self
3006 .handle_list_clap_plugins_with_capabilities(a.clone())
3007 .await
3008 {
3009 return;
3010 }
3011 }
3012 Action::ClapPlugins(_) => {}
3013 Action::ClapPluginsUnavailable { .. } => {}
3014 Action::TrackLoadClapPlugin {
3015 ref track_name,
3016 ref plugin_id,
3017 instance_id,
3018 } => {
3019 if self
3020 .handle_track_load_clap_plugin(
3021 track_name.as_str(),
3022 plugin_id.as_str(),
3023 instance_id,
3024 )
3025 .await
3026 {
3027 return;
3028 }
3029 }
3030 Action::TrackUnloadClapPlugin {
3031 ref track_name,
3032 ref plugin_id,
3033 } => {
3034 if self
3035 .handle_track_unload_clap_plugin(track_name.as_str(), plugin_id.as_str())
3036 .await
3037 {
3038 return;
3039 }
3040 }
3041 Action::TrackUnloadClapPluginInstance {
3042 ref track_name,
3043 instance_id,
3044 } => {
3045 if self
3046 .handle_track_unload_clap_plugin_instance(track_name.as_str(), instance_id)
3047 .await
3048 {
3049 return;
3050 }
3051 }
3052 Action::TrackShowClapGui { .. } => {
3053 if Self::box_bool(self.handle_track_show_clap_gui(a.clone())).await {
3054 return;
3055 }
3056 }
3057 Action::TrackLoadVst3Plugin {
3058 ref track_name,
3059 ref plugin_id,
3060 instance_id,
3061 } => {
3062 if self
3063 .handle_track_load_vst3_plugin(
3064 track_name.as_str(),
3065 plugin_id.as_str(),
3066 instance_id,
3067 )
3068 .await
3069 {
3070 return;
3071 }
3072 }
3073 Action::TrackUnloadVst3Plugin {
3074 ref track_name,
3075 ref plugin_id,
3076 } => {
3077 if self
3078 .handle_track_unload_vst3_plugin(track_name.as_str(), plugin_id.as_str())
3079 .await
3080 {
3081 return;
3082 }
3083 }
3084 Action::TrackUnloadVst3PluginInstance {
3085 ref track_name,
3086 instance_id,
3087 } => {
3088 if self
3089 .handle_track_unload_vst3_plugin_instance(track_name.as_str(), instance_id)
3090 .await
3091 {
3092 return;
3093 }
3094 }
3095 Action::TrackShowVst3Gui { .. } => {
3096 if Self::box_bool(self.handle_track_show_vst3_gui(a.clone())).await {
3097 return;
3098 }
3099 }
3100 #[cfg(all(unix, not(target_os = "macos")))]
3101 Action::TrackLoadLv2Plugin {
3102 ref track_name,
3103 ref plugin_uri,
3104 instance_id,
3105 } => {
3106 if self
3107 .handle_track_load_lv2_plugin(
3108 track_name.as_str(),
3109 plugin_uri.as_str(),
3110 instance_id,
3111 )
3112 .await
3113 {
3114 return;
3115 }
3116 }
3117 #[cfg(all(unix, not(target_os = "macos")))]
3118 Action::TrackUnloadLv2Plugin {
3119 ref track_name,
3120 ref plugin_uri,
3121 } => {
3122 if self
3123 .handle_track_unload_lv2_plugin(track_name.as_str(), plugin_uri.as_str())
3124 .await
3125 {
3126 return;
3127 }
3128 }
3129 #[cfg(all(unix, not(target_os = "macos")))]
3130 Action::TrackUnloadLv2PluginInstance {
3131 ref track_name,
3132 instance_id,
3133 } => {
3134 if self
3135 .handle_track_unload_lv2_plugin_instance(track_name.as_str(), instance_id)
3136 .await
3137 {
3138 return;
3139 }
3140 }
3141 #[cfg(all(unix, not(target_os = "macos")))]
3142 Action::TrackShowLv2Gui { .. } => {
3143 if Self::box_bool(self.handle_track_show_lv2_gui(a.clone())).await {
3144 return;
3145 }
3146 }
3147 Action::TrackSetPluginResourceDir { .. } => {
3148 if Self::box_bool(self.handle_track_set_plugin_resource_dir(a.clone())).await {
3149 return;
3150 }
3151 }
3152 Action::TrackClapFileReferences { .. } => {
3153 if Self::box_bool(self.handle_track_clap_file_references(a.clone())).await {
3154 return;
3155 }
3156 }
3157 Action::TrackUpdateClapFileReference { .. } => {
3158 if self
3159 .handle_track_update_clap_file_reference(a.clone())
3160 .await
3161 {
3162 return;
3163 }
3164 }
3165 Action::ClipSetPluginResourceDir { .. } => {
3166 if Self::box_bool(self.handle_clip_set_plugin_resource_dir(a.clone())).await {
3167 return;
3168 }
3169 }
3170 Action::ClipClapFileReferences { .. } => {
3171 if Self::box_bool(self.handle_clip_clap_file_references(a.clone())).await {
3172 return;
3173 }
3174 }
3175 Action::ClipUpdateClapFileReference { .. } => {
3176 if Self::box_bool(self.handle_clip_update_clap_file_reference(a.clone())).await {
3177 return;
3178 }
3179 }
3180 Action::TrackSetClapParameter { .. } => {
3181 if Self::box_bool(self.handle_track_set_clap_parameter(a.clone())).await {
3182 return;
3183 }
3184 }
3185 Action::ClipSetClapParameter { .. } => {
3186 if Self::box_bool(self.handle_clip_set_clap_parameter(a.clone())).await {
3187 return;
3188 }
3189 }
3190 Action::TrackSetClapParameterAt { .. } => {
3191 if Self::box_bool(self.handle_track_set_clap_parameter_at(a.clone())).await {
3192 return;
3193 }
3194 }
3195 Action::TrackBeginClapParameterEdit { .. } => {
3196 if Self::box_bool(self.handle_track_begin_clap_parameter_edit(a.clone())).await {
3197 return;
3198 }
3199 }
3200 Action::TrackEndClapParameterEdit { .. } => {
3201 if Self::box_bool(self.handle_track_end_clap_parameter_edit(a.clone())).await {
3202 return;
3203 }
3204 }
3205 Action::TrackGetClapParameters { .. } => {
3206 if Self::box_bool(self.handle_track_get_clap_parameters(a.clone())).await {
3207 return;
3208 }
3209 }
3210 Action::TrackClapParameters { .. } => {}
3211 Action::TrackClapSnapshotState { .. } => {
3212 if Self::box_bool(self.handle_track_clap_snapshot_state(a.clone())).await {
3213 return;
3214 }
3215 }
3216 Action::ClipClapSnapshotState { .. } => {
3217 if Self::box_bool(self.handle_clip_clap_snapshot_state(a.clone())).await {
3218 return;
3219 }
3220 }
3221 Action::TrackClapStateSnapshot { .. } => {}
3222 Action::ClipClapStateSnapshot { .. } => {}
3223 Action::TrackClapStateDirty { .. } => {}
3224 Action::ClipClapStateDirty { .. } => {}
3225 Action::TrackClapRestoreState { .. } => {
3226 if Self::box_bool(self.handle_track_clap_restore_state(a.clone())).await {
3227 return;
3228 }
3229 }
3230 Action::ClipClapRestoreState { .. } => {
3231 if Self::box_bool(self.handle_clip_clap_restore_state(a.clone())).await {
3232 return;
3233 }
3234 }
3235 Action::TrackSnapshotAllClapStates { .. } => {
3236 if Self::box_bool(self.handle_track_snapshot_all_clap_states(a.clone())).await {
3237 return;
3238 }
3239 }
3240 Action::TrackSnapshotAllClapStatesDone { .. } => {}
3241 Action::TrackGetVst3Graph { .. } => {
3242 if Self::box_bool(self.handle_track_get_vst3_graph(a.clone())).await {
3243 return;
3244 }
3245 }
3246 Action::TrackVst3Graph { .. } => {}
3247 Action::TrackSetVst3Parameter { .. } => {
3248 if Self::box_bool(self.handle_track_set_vst3_parameter(a.clone())).await {
3249 return;
3250 }
3251 }
3252 Action::TrackSetPluginBypassed { .. } => {
3253 if Self::box_bool(self.handle_track_set_plugin_bypassed(a.clone())).await {
3254 return;
3255 }
3256 }
3257 Action::TrackGetVst3Parameters { .. } => {
3258 if Self::box_bool(self.handle_track_get_vst3_parameters(a.clone())).await {
3259 return;
3260 }
3261 }
3262 Action::TrackVst3Parameters { .. } => {}
3263 #[cfg(all(unix, not(target_os = "macos")))]
3264 Action::TrackGetLv2PluginControls { .. } => {
3265 if Self::box_bool(self.handle_track_get_lv2_plugin_controls(a.clone())).await {
3266 return;
3267 }
3268 }
3269 #[cfg(all(unix, not(target_os = "macos")))]
3270 Action::TrackLv2SnapshotState { .. } => {
3271 if Self::box_bool(self.handle_track_lv2_snapshot_state(a.clone())).await {
3272 return;
3273 }
3274 }
3275 #[cfg(all(unix, not(target_os = "macos")))]
3276 Action::ClipLv2SnapshotState { .. } => {
3277 if Self::box_bool(self.handle_clip_lv2_snapshot_state(a.clone())).await {
3278 return;
3279 }
3280 }
3281 Action::TrackVst3SnapshotState { .. } => {
3282 if Self::box_bool(self.handle_track_vst3_snapshot_state(a.clone())).await {
3283 return;
3284 }
3285 }
3286 Action::ClipVst3SnapshotState { .. } => {
3287 if Self::box_bool(self.handle_clip_vst3_snapshot_state(a.clone())).await {
3288 return;
3289 }
3290 }
3291 Action::TrackVst3StateSnapshot { .. } => {}
3292 Action::ClipVst3StateSnapshot { .. } => {}
3293 Action::TrackVst3RestoreState { .. } => {
3294 if Self::box_bool(self.handle_track_vst3_restore_state(a.clone())).await {
3295 return;
3296 }
3297 }
3298 Action::TrackConnectVst3Audio { .. } => {
3299 if Self::box_bool(self.handle_track_connect_vst3_audio(a.clone())).await {
3300 return;
3301 }
3302 }
3303 Action::TrackDisconnectVst3Audio { .. } => {
3304 if Self::box_bool(self.handle_track_disconnect_vst3_audio(a.clone())).await {
3305 return;
3306 }
3307 }
3308 Action::ClipMove { .. } => {
3309 self.handle_clip_move(a.clone()).await;
3310 }
3311 Action::AddClip { .. } => {
3312 if Self::box_bool(self.handle_add_clip(a.clone())).await {
3313 return;
3314 }
3315 }
3316 Action::AddGroupedClip { .. } => {
3317 if Self::box_bool(self.handle_add_grouped_clip(a.clone())).await {
3318 return;
3319 }
3320 }
3321 Action::RemoveClip {
3322 ref track_name,
3323 kind,
3324 ref clip_indices,
3325 } => {
3326 self.remove_clips_from_track(track_name, kind, clip_indices);
3327 }
3328 Action::MoveClipToUnused {
3329 ref track_name,
3330 kind,
3331 ref clip_indices,
3332 } => {
3333 self.move_clips_to_unused(track_name, kind, clip_indices);
3334 }
3335 Action::DeleteUnusedClips { ref clip_ids } => {
3336 self.delete_unused_clips(clip_ids);
3337 }
3338 Action::SetUnusedClips {
3339 ref audio,
3340 ref midi,
3341 } => {
3342 self.set_unused_clips(audio.clone(), midi.clone());
3343 }
3344 Action::RenameClip {
3345 ref track_name,
3346 kind,
3347 clip_index,
3348 ref new_name,
3349 } => {
3350 self.rename_clip_references(track_name, kind, clip_index, new_name);
3351 }
3352 Action::SetClipIdentity {
3353 ref track_name,
3354 kind,
3355 clip_index,
3356 ref new_id,
3357 ref new_name,
3358 } => {
3359 self.set_clip_identity(track_name, kind, clip_index, new_id, new_name);
3360 }
3361 Action::SetClipSourceName {
3362 ref track_name,
3363 kind,
3364 clip_index,
3365 ref name,
3366 } => {
3367 self.set_clip_source_name(track_name, clip_index, kind, name.clone());
3368 }
3369 Action::SetClipFade { .. } => {
3370 if Self::box_bool(self.handle_set_clip_fade(a.clone())).await {
3371 return;
3372 }
3373 }
3374 Action::SetClipBounds {
3375 ref track_name,
3376 clip_index,
3377 kind,
3378 start,
3379 length,
3380 offset,
3381 } => {
3382 self.set_clip_bounds(track_name, clip_index, kind, start, length, offset);
3383 }
3384 Action::SyncClipBounds {
3385 ref track_name,
3386 clip_index,
3387 kind,
3388 start,
3389 length,
3390 offset,
3391 } => {
3392 self.set_clip_bounds(track_name, clip_index, kind, start, length, offset);
3393 }
3394 Action::SetClipMuted {
3395 ref track_name,
3396 clip_index,
3397 kind,
3398 muted,
3399 } => {
3400 self.set_clip_muted(track_name, clip_index, kind, muted);
3401 }
3402 Action::SetClipPluginGraphJson {
3403 ref track_name,
3404 clip_index,
3405 ref plugin_graph_json,
3406 } => {
3407 self.set_clip_plugin_graph_json(track_name, clip_index, plugin_graph_json.clone());
3408 }
3409 Action::SetClipPitchCorrection { .. } => {
3410 if Self::box_bool(self.handle_set_clip_pitch_correction(a.clone())).await {
3411 return;
3412 }
3413 }
3414 Action::Connect {
3415 ref from_track,
3416 from_port,
3417 ref to_track,
3418 to_port,
3419 kind,
3420 } => {
3421 self.handle_connect(
3422 from_track.as_str(),
3423 from_port,
3424 to_track.as_str(),
3425 to_port,
3426 kind,
3427 )
3428 .await;
3429 }
3430 Action::Disconnect { .. } => {
3431 self.handle_disconnect(a.clone()).await;
3432 }
3433 Action::OpenAudioDevice { .. } => {
3434 let (done, updated) = self.handle_open_audio_device(a.clone()).await;
3435 if done {
3436 return;
3437 }
3438 if let Some(action) = updated {
3439 action_to_process = action;
3440 }
3441 }
3442 Action::JackAddAudioInputPort => {
3443 if Self::box_bool(self.handle_jack_add_audio_input_port(a.clone())).await {
3444 return;
3445 }
3446 }
3447 Action::JackRemoveAudioInputPort(_removed_port) => {
3448 if self
3449 .handle_jack_remove_audio_input_port(_removed_port, a.clone())
3450 .await
3451 {
3452 return;
3453 }
3454 }
3455 Action::JackAddAudioOutputPort => {
3456 if Self::box_bool(self.handle_jack_add_audio_output_port(a.clone())).await {
3457 return;
3458 }
3459 }
3460 Action::JackRemoveAudioOutputPort(_removed_port) => {
3461 if self
3462 .handle_jack_remove_audio_output_port(_removed_port, a.clone())
3463 .await
3464 {
3465 return;
3466 }
3467 }
3468 Action::JackGetGraph => {
3469 #[cfg(unix)]
3470 {
3471 match self
3472 .jack_runtime
3473 .as_ref()
3474 .ok_or(
3475 "JACK runtime is not active; open the JACK backend first".to_string(),
3476 )
3477 .and_then(|jack| jack.graph_info())
3478 {
3479 Ok(graph) => self.notify_clients(Ok(Action::JackGraph(graph))).await,
3480 Err(e) => self.notify_clients(Err(e)).await,
3481 }
3482 }
3483 #[cfg(not(unix))]
3484 {
3485 self.notify_clients(Err(
3486 "JACK backend is not available on this platform build".to_string(),
3487 ))
3488 .await;
3489 }
3490 return;
3491 }
3492 Action::JackConnect {
3493 ref source,
3494 ref destination,
3495 } => {
3496 #[cfg(unix)]
3497 {
3498 match self
3499 .jack_runtime
3500 .as_ref()
3501 .ok_or(
3502 "JACK runtime is not active; open the JACK backend first".to_string(),
3503 )
3504 .and_then(|jack| jack.connect_ports_by_name(source, destination))
3505 .and_then(|_| {
3506 self.jack_runtime
3507 .as_ref()
3508 .expect("JACK runtime was checked")
3509 .graph_info()
3510 }) {
3511 Ok(graph) => {
3512 self.notify_clients(Ok(a.clone())).await;
3513 self.notify_clients(Ok(Action::JackGraph(graph))).await;
3514 }
3515 Err(e) => self.notify_clients(Err(e)).await,
3516 }
3517 }
3518 #[cfg(not(unix))]
3519 {
3520 self.notify_clients(Err(
3521 "JACK backend is not available on this platform build".to_string(),
3522 ))
3523 .await;
3524 }
3525 return;
3526 }
3527 Action::JackDisconnect {
3528 ref source,
3529 ref destination,
3530 } => {
3531 #[cfg(unix)]
3532 {
3533 match self
3534 .jack_runtime
3535 .as_ref()
3536 .ok_or(
3537 "JACK runtime is not active; open the JACK backend first".to_string(),
3538 )
3539 .and_then(|jack| jack.disconnect_ports_by_name(source, destination))
3540 .and_then(|_| {
3541 self.jack_runtime
3542 .as_ref()
3543 .expect("JACK runtime was checked")
3544 .graph_info()
3545 }) {
3546 Ok(graph) => {
3547 self.notify_clients(Ok(a.clone())).await;
3548 self.notify_clients(Ok(Action::JackGraph(graph))).await;
3549 }
3550 Err(e) => self.notify_clients(Err(e)).await,
3551 }
3552 }
3553 #[cfg(not(unix))]
3554 {
3555 self.notify_clients(Err(
3556 "JACK backend is not available on this platform build".to_string(),
3557 ))
3558 .await;
3559 }
3560 return;
3561 }
3562 Action::OpenMidiInputDevice(ref device) => {
3563 if let Some(worker) = &self.hw_worker {
3564 if let Err(e) = worker
3565 .tx
3566 .send(Message::HWOpenMidiInputDevice(device.clone()))
3567 .await
3568 {
3569 self.notify_clients(Err(format!("Failed to send MIDI input open: {e}")))
3570 .await;
3571 }
3572 return;
3573 }
3574 let Some(midi_hub) = self.midi_hub.as_mut() else {
3575 self.notify_clients(Err("Hardware MIDI hub is not available".to_string()))
3576 .await;
3577 return;
3578 };
3579 if let Err(e) = midi_hub.open_input(device) {
3580 self.notify_clients(Err(e)).await;
3581 return;
3582 }
3583 }
3584 Action::OpenMidiOutputDevice(ref device) => {
3585 if let Some(worker) = &self.hw_worker {
3586 if let Err(e) = worker
3587 .tx
3588 .send(Message::HWOpenMidiOutputDevice(device.clone()))
3589 .await
3590 {
3591 self.notify_clients(Err(format!("Failed to send MIDI output open: {e}")))
3592 .await;
3593 }
3594 return;
3595 }
3596 let Some(midi_hub) = self.midi_hub.as_mut() else {
3597 self.notify_clients(Err("Hardware MIDI hub is not available".to_string()))
3598 .await;
3599 return;
3600 };
3601 if let Err(e) = midi_hub.open_output(device) {
3602 self.notify_clients(Err(e)).await;
3603 return;
3604 }
3605 }
3606 Action::RequestSessionDiagnostics => {
3607 self.handle_request_session_diagnostics().await;
3608 }
3609 Action::RequestMidiLearnMappingsReport => {
3610 self.handle_request_midi_learn_mappings_report().await;
3611 }
3612 Action::ClearAllMidiLearnBindings => {
3613 if Self::box_bool(self.handle_clear_all_midi_learn_bindings(a.clone())).await {
3614 return;
3615 }
3616 }
3617 #[cfg(all(unix, not(target_os = "macos")))]
3618 Action::TrackLv2PluginControls { .. } => {}
3619 #[cfg(all(unix, not(target_os = "macos")))]
3620 Action::ClipLv2PluginControls { .. } => {}
3621 #[cfg(all(unix, not(target_os = "macos")))]
3622 Action::TrackLv2StateSnapshot { .. } => {}
3623 #[cfg(all(unix, not(target_os = "macos")))]
3624 Action::ClipLv2StateSnapshot { .. } => {}
3625 #[cfg(all(unix, not(target_os = "macos")))]
3626 Action::TrackLv2Midnam { .. } => {}
3627 Action::TrackClapNoteNames { .. } => {}
3628 Action::SessionDiagnosticsReport { .. } => {}
3629 Action::MidiLearnMappingsReport { .. } => {}
3630 Action::HWInfo { .. } => {}
3631 Action::HistoryState { .. } => {}
3632 Action::Undo => {}
3633 Action::Redo => {}
3634 Action::ApplyGroupedActions(_) => {}
3635 _ => {}
3636 }
3637
3638 if let Some(inverse) = inverse_actions {
3639 if let Some(group) = self.history_group.as_mut() {
3640 group.forward_actions.push(action_to_process.clone());
3641 group.inverse_actions.splice(0..0, inverse);
3642 } else {
3643 self.history.record(UndoEntry {
3644 forward_actions: vec![action_to_process.clone()],
3645 inverse_actions: inverse,
3646 });
3647 }
3648 }
3649
3650 self.notify_clients(Ok(action_to_process)).await;
3651 }
3652 pub async fn work(&mut self) {
3653 let mut tick = tokio::time::interval(Duration::from_millis(1));
3656 tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
3657 loop {
3658 let message = tokio::select! {
3659 message = self.rx.recv() => {
3660 let Some(message) = message else {
3661 break;
3662 };
3663 message
3664 }
3665 _ = tick.tick() => {
3666 self.poll_node_worker_results().await;
3667 self.poll_jack_hw_finished().await;
3668 self.on_executor_tick().await;
3669 continue;
3670 }
3671 };
3672 self.poll_node_worker_results().await;
3673 self.poll_jack_hw_finished().await;
3674 match message {
3675 Message::Ready(id) => {
3676 if let Some(track_name) = self.bounce_worker_tracks.remove(&id) {
3680 self.offline_bounce_jobs.remove(&track_name);
3681 }
3682 self.push_ready_worker(id);
3683 if self.dispatch_node_jobs(Vec::new()).await {
3684 self.on_all_tracks_finished().await;
3685 }
3686 self.drain_pending_requests_if_idle().await;
3687 }
3688 Message::NodeDone {
3689 worker_id,
3690 epoch,
3691 node,
3692 output_linear,
3693 parameter_updates,
3694 latency_changed,
3695 } => {
3696 self.on_node_done(
3697 worker_id,
3698 epoch,
3699 node,
3700 output_linear,
3701 parameter_updates,
3702 latency_changed,
3703 )
3704 .await;
3705 }
3706 Message::Channel(s) => {
3707 self.clients.push(s);
3708 }
3709 Message::Response(result) => {
3710 self.notify_clients(result).await;
3711 }
3712
3713 Message::Request(a) => {
3714 self.dispatch_request(a).await;
3715 self.plan_builder.mark_dirty();
3718 }
3719 Message::OscRequest { action, reply_to } => {
3720 self.osc_reply_target = Some(reply_to);
3721 self.dispatch_request(action).await;
3722 self.osc_reply_target = None;
3723 self.plan_builder.mark_dirty();
3724 }
3725 Message::OfflineBounceFinished { result } => {
3726 if let Ok(Action::TrackOfflineBounce { track_name, .. })
3727 | Ok(Action::TrackOfflineBounceCanceled { track_name, .. }) = &result
3728 {
3729 self.offline_bounce_jobs.remove(track_name);
3730 }
3731 self.notify_clients(result).await;
3732 self.drain_pending_requests_if_idle().await;
3733 }
3734 Message::HWFinished => {
3735 if !self.awaiting_hwfinished {
3736 tracing::debug!("HWFinished ignored (not awaiting)");
3737 continue;
3738 }
3739 tracing::debug!("HWFinished handling; playing={}", self.playing);
3740 self.handling_hwfinished = true;
3741 self.awaiting_hwfinished = false;
3742 #[cfg(unix)]
3743 {
3744 if let Some(jack) = self.jack_runtime.as_mut() {
3745 if !self.pending_hw_midi_out_events.is_empty() {
3746 let out_events =
3747 std::mem::take(&mut self.pending_hw_midi_out_events);
3748 jack.write_events(&out_events);
3749 }
3750 let mut in_events = vec![];
3751 jack.read_events_into(&mut in_events);
3752 if !in_events.is_empty() {
3753 self.pending_hw_midi_events.extend(in_events);
3754 }
3755 let dropped = jack.take_midi_events_dropped();
3756 if dropped > 0 {
3757 tracing::warn!(
3758 "JACK MIDI ring full; {dropped} events dropped since last cycle"
3759 );
3760 }
3761 }
3762 }
3763 #[cfg(unix)]
3764 if self.jack_runtime.is_some() {
3765 self.sync_from_jack_transport().await;
3766 }
3767 while let Some(a) = self.pending_requests.pop_front() {
3768 self.handle_request(a).await;
3769 }
3770 self.apply_mute_solo_policy();
3771 self.append_recorded_cycle();
3772 self.flush_completed_recordings().await;
3773 let hw_in_routes = self.midi_hw_in_routes.clone();
3774 let pending_hw_in_by_device = self.pending_hw_midi_events_by_device.clone();
3775 let mut reconfigured_tracks = Vec::new();
3776 let state = self.state_snapshot.load_full();
3777 for (track_name, track) in state.tracks.iter() {
3778 let mut track_lock = track.lock();
3779 if self.jack_runtime_is_some() {
3780 if !self.pending_hw_midi_events.is_empty() {
3781 track_lock.push_hw_midi_events(&self.pending_hw_midi_events);
3782 }
3783 } else {
3784 for route in hw_in_routes.iter().filter(|r| &r.to_track == track_name) {
3785 if let Some(events) = pending_hw_in_by_device.get(&route.device) {
3786 track_lock.push_hw_midi_events_to_port(route.to_port, events);
3787 }
3788 }
3789 }
3790 if track_lock.setup() {
3791 reconfigured_tracks.push(track_name.clone());
3792 }
3793 }
3794 self.publish_track_meters();
3795 self.publish_session_runtime_reports().await;
3796 self.publish_clap_state_dirty().await;
3797 for track_name in reconfigured_tracks {
3798 let track = state.tracks.get(&track_name).cloned();
3799 if let Some(track) = track {
3800 let (plugins, connections, connectable_connections) = {
3801 let track_lock = track.lock();
3802 (
3803 track_lock.plugin_graph_plugins(false),
3804 track_lock.plugin_graph_connections(),
3805 track_lock.connectable_connections(),
3806 )
3807 };
3808 self.notify_clients(Ok(Action::TrackPluginGraph {
3809 track_name: track_name.clone(),
3810 plugins,
3811 connections,
3812 connectable_connections,
3813 }))
3814 .await;
3815 }
3816 }
3817 self.pending_hw_midi_events.clear();
3818 self.pending_hw_midi_events_by_device.clear();
3819 if self.transport_running {
3820 if self.transport_panic_flush_pending {
3821 self.transport_panic_flush_pending = false;
3822 } else if self.transport_restart_pending {
3823 self.transport_restart_pending = false;
3824 } else {
3825 let next = self
3826 .transport_sample
3827 .saturating_add(self.current_cycle_samples());
3828 let normalized = self.normalize_transport_sample(next);
3829 let wrapped = normalized != next;
3830 self.transport_sample = normalized;
3831 self.publish_transport_snapshot();
3832 if wrapped {
3833 if self.notified_loop_wrap_sample == Some(self.transport_sample) {
3834 self.notified_loop_wrap_sample = None;
3835 } else {
3836 self.notify_clients(Ok(Action::TransportPosition(
3837 self.transport_sample,
3838 )))
3839 .await;
3840 }
3841 }
3842 }
3843 }
3844 if self.session_clip_playback_enabled && self.playing {
3845 self.session_transport_sample = self
3846 .session_transport_sample
3847 .saturating_add(self.current_cycle_samples());
3848 }
3849 {
3850 let echoes = self.apply_modulators(self.active_transport_sample());
3851 for action in echoes {
3852 self.notify_clients(Ok(action)).await;
3853 }
3854 }
3855 self.start_plan_cycle().await;
3856 #[cfg(unix)]
3857 {
3858 if self.jack_runtime.is_some() {
3859 self.awaiting_hwfinished = true;
3860 }
3861 }
3862 self.handling_hwfinished = false;
3863 }
3864 Message::HWMidiEvents(events) => {
3865 for hw_event in events {
3866 let thru_targets: Vec<String> = self
3867 .midi_hw_thru_routes
3868 .iter()
3869 .filter(|route| route.from_device == hw_event.device)
3870 .map(|route| route.to_device.clone())
3871 .collect();
3872 for device in thru_targets {
3873 self.pending_hw_midi_out_events_by_device.push(HwMidiEvent {
3874 device,
3875 event: hw_event.event.clone(),
3876 });
3877 }
3878 if hw_event.event.data.len() >= 3 {
3879 let status = hw_event.event.data[0];
3880 if status & 0xF0 == 0xB0 {
3881 let channel = status & 0x0F;
3882 let cc = hw_event.event.data[1];
3883 let value = hw_event.event.data[2];
3884 self.handle_incoming_hw_cc(&hw_event.device, channel, cc, value)
3885 .await;
3886 }
3887 if self.step_recording_enabled && status & 0xF0 == 0x90 {
3888 let channel = status & 0x0F;
3889 let pitch = hw_event.event.data[1];
3890 let velocity = hw_event.event.data[2];
3891 if velocity > 0 {
3892 self.notify_clients(Ok(Action::StepRecordMidiNote {
3893 device: hw_event.device.clone(),
3894 channel,
3895 pitch,
3896 velocity,
3897 }))
3898 .await;
3899 }
3900 }
3901 }
3902 self.pending_hw_midi_events_by_device
3903 .entry(hw_event.device)
3904 .or_default()
3905 .push(hw_event.event);
3906 }
3907 }
3908 _ => {}
3909 }
3910 }
3911 }
3912
3913 pub(crate) fn collect_hw_midi_output_events(&self) -> Vec<MidiEvent> {
3914 let mut events = vec![];
3915 for track in self.state_snapshot.load_full().tracks.values() {
3916 events.extend(
3917 track
3918 .lock()
3919 .take_hw_midi_out_events()
3920 .into_iter()
3921 .map(|evt| evt.event),
3922 );
3923 }
3924 events.sort_by_key(|a| a.frame);
3925 events
3926 }
3927
3928 pub(crate) fn collect_hw_midi_output_events_by_device(&mut self) -> Vec<HwMidiEvent> {
3929 let mut events = Vec::<HwMidiEvent>::new();
3930 let routes = self.midi_hw_out_routes.clone();
3931 let mut events_by_track = HashMap::<String, Vec<crate::track::HwMidiOutEvent>>::new();
3932 {
3933 let state = self.state_snapshot.load_full();
3934 for route in &routes {
3935 if events_by_track.contains_key(&route.from_track) {
3936 continue;
3937 }
3938 let Some(track) = state.tracks.get(&route.from_track) else {
3939 continue;
3940 };
3941 events_by_track.insert(
3942 route.from_track.clone(),
3943 track.lock().take_hw_midi_out_events(),
3944 );
3945 }
3946 }
3947
3948 for route in routes {
3949 let Some(track_events) = events_by_track.get(&route.from_track) else {
3950 continue;
3951 };
3952 for hw_event in track_events
3953 .iter()
3954 .filter(|evt| evt.port == route.from_port)
3955 {
3956 self.update_active_hw_notes_for_track(
3957 &route.from_track,
3958 &route.device,
3959 &hw_event.event.data,
3960 );
3961 events.push(HwMidiEvent {
3962 device: route.device.clone(),
3963 event: hw_event.event.clone(),
3964 });
3965 }
3966 }
3967 events.sort_by(|a, b| {
3968 a.event
3969 .frame
3970 .cmp(&b.event.frame)
3971 .then_with(|| a.device.cmp(&b.device))
3972 });
3973 events
3974 }
3975}