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_stopped_plugin_parameter_echoes(&mut self) {
1577 if self.playing || self.transport_running {
1578 return;
1579 }
1580
1581 let state = self.state_snapshot.load_full();
1582 let mut updates = Vec::new();
1583 for track in state.tracks.values() {
1584 updates.extend(track.lock().drain_plugin_parameter_echoes());
1585 }
1586 drop(state);
1587
1588 for action in updates {
1589 self.notify_clients(Ok(action)).await;
1590 }
1591 }
1592
1593 pub(crate) async fn poll_node_worker_results(&mut self) {
1594 let mut results = Vec::new();
1595 for worker in &mut self.workers {
1596 if let Some(rx) = worker.node_result_rx.as_mut() {
1597 while let Ok(result) = rx.pop() {
1598 results.push(result);
1599 }
1600 }
1601 }
1602 for result in results {
1603 self.on_node_done(
1604 result.worker_id,
1605 result.epoch,
1606 result.node,
1607 result.output_linear,
1608 result.parameter_updates,
1609 result.latency_changed,
1610 )
1611 .await;
1612 }
1613 }
1614
1615 pub(crate) async fn poll_jack_hw_finished(&mut self) {
1616 #[cfg(unix)]
1617 {
1618 let finished = self
1619 .jack_runtime
1620 .as_ref()
1621 .map(|jack| jack.take_hw_finished_count())
1622 .unwrap_or(0);
1623 if finished > 0 {
1624 self.handle_hw_finished().await;
1625 }
1626 }
1627 }
1628
1629 pub(crate) async fn handle_hw_finished(&mut self) {
1630 if !self.awaiting_hwfinished {
1631 tracing::debug!("HWFinished ignored (not awaiting)");
1632 return;
1633 }
1634 tracing::debug!("HWFinished handling; playing={}", self.playing);
1635 self.handling_hwfinished = true;
1636 self.awaiting_hwfinished = false;
1637 #[cfg(unix)]
1638 {
1639 if let Some(jack) = self.jack_runtime.as_mut() {
1640 if !self.pending_hw_midi_out_events.is_empty() {
1641 let out_events = std::mem::take(&mut self.pending_hw_midi_out_events);
1642 jack.write_events(&out_events);
1643 }
1644 let mut in_events = vec![];
1645 jack.read_events_into(&mut in_events);
1646 if !in_events.is_empty() {
1647 self.pending_hw_midi_events.extend(in_events);
1648 }
1649 let dropped = jack.take_midi_events_dropped();
1650 if dropped > 0 {
1651 tracing::warn!(
1652 "JACK MIDI ring full; {dropped} events dropped since last cycle"
1653 );
1654 }
1655 }
1656 }
1657 #[cfg(unix)]
1658 if self.jack_runtime.is_some() {
1659 self.sync_from_jack_transport().await;
1660 }
1661 while let Some(a) = self.pending_requests.pop_front() {
1662 self.handle_request(a).await;
1663 }
1664 self.apply_mute_solo_policy();
1665 self.append_recorded_cycle();
1666 self.flush_completed_recordings().await;
1667 let hw_in_routes = self.midi_hw_in_routes.clone();
1668 let pending_hw_in_by_device = self.pending_hw_midi_events_by_device.clone();
1669 let mut reconfigured_tracks = Vec::new();
1670 let state = self.state_snapshot.load_full();
1671 for (track_name, track) in state.tracks.iter() {
1672 let mut track_lock = track.lock();
1673 if self.jack_runtime_is_some() {
1674 if !self.pending_hw_midi_events.is_empty() {
1675 track_lock.push_hw_midi_events(&self.pending_hw_midi_events);
1676 }
1677 } else {
1678 for route in hw_in_routes.iter().filter(|r| &r.to_track == track_name) {
1679 if let Some(events) = pending_hw_in_by_device.get(&route.device) {
1680 track_lock.push_hw_midi_events_to_port(route.to_port, events);
1681 }
1682 }
1683 }
1684 if track_lock.setup() {
1685 reconfigured_tracks.push(track_name.clone());
1686 }
1687 }
1688 self.publish_track_meters();
1689 self.publish_session_runtime_reports().await;
1690 self.publish_clap_state_dirty().await;
1691 for track_name in reconfigured_tracks {
1692 let track = state.tracks.get(&track_name).cloned();
1693 if let Some(track) = track {
1694 let (plugins, connections, connectable_connections) = {
1695 let track_lock = track.lock();
1696 (
1697 track_lock.plugin_graph_plugins(false),
1698 track_lock.plugin_graph_connections(),
1699 track_lock.connectable_connections(),
1700 )
1701 };
1702 self.notify_clients(Ok(Action::TrackPluginGraph {
1703 track_name: track_name.clone(),
1704 plugins,
1705 connections,
1706 connectable_connections,
1707 }))
1708 .await;
1709 }
1710 }
1711 self.pending_hw_midi_events.clear();
1712 self.pending_hw_midi_events_by_device.clear();
1713 if self.transport_running {
1714 if self.transport_panic_flush_pending {
1715 self.transport_panic_flush_pending = false;
1716 } else if self.transport_restart_pending {
1717 self.transport_restart_pending = false;
1718 } else {
1719 let next = self
1720 .transport_sample
1721 .saturating_add(self.current_cycle_samples());
1722 let normalized = self.normalize_transport_sample(next);
1723 let wrapped = normalized != next;
1724 self.transport_sample = normalized;
1725 self.publish_transport_snapshot();
1726 if wrapped {
1727 if self.notified_loop_wrap_sample == Some(self.transport_sample) {
1728 self.notified_loop_wrap_sample = None;
1729 } else {
1730 self.notify_clients(Ok(Action::TransportPosition(self.transport_sample)))
1731 .await;
1732 }
1733 }
1734 }
1735 }
1736 if self.session_clip_playback_enabled && self.playing {
1737 self.session_transport_sample = self
1738 .session_transport_sample
1739 .saturating_add(self.current_cycle_samples());
1740 }
1741 {
1742 let echoes = self.apply_modulators(self.active_transport_sample());
1743 for action in echoes {
1744 self.notify_clients(Ok(action)).await;
1745 }
1746 }
1747 self.start_plan_cycle().await;
1748 #[cfg(unix)]
1749 {
1750 if self.jack_runtime.is_some() {
1751 self.awaiting_hwfinished = true;
1752 }
1753 }
1754 self.handling_hwfinished = false;
1755 }
1756
1757 pub(crate) async fn on_executor_tick(&mut self) {
1760 if self.executor.cycle_complete() {
1761 return;
1762 }
1763 let outcome = self
1764 .executor
1765 .force_timeouts(Instant::now(), Self::TRACK_PROCESS_TIMEOUT);
1766 self.log_silenced_nodes(&outcome.silenced);
1767 let mut complete = self.dispatch_node_jobs(outcome.jobs).await;
1768 complete |= outcome.cycle_complete;
1769 if complete {
1770 self.on_all_tracks_finished().await;
1771 }
1772 }
1773
1774 pub(crate) fn log_silenced_nodes(&self, silenced: &[u32]) {
1775 for &node in silenced {
1776 let plan = self.executor.plan();
1777 let name = match plan.nodes.get(node as usize) {
1778 Some(crate::render_plan::Op::Task { task, .. }) => Self::task_track_name(task),
1779 _ => format!("node {node}"),
1780 };
1781 tracing::warn!(
1782 "Node {} ('{}') exceeded process timeout ({} ms); forced silent completion for cycle",
1783 node,
1784 name,
1785 Self::TRACK_PROCESS_TIMEOUT.as_millis()
1786 );
1787 }
1788 }
1789
1790 pub(crate) async fn send_bounce_job(
1793 &mut self,
1794 worker_index: usize,
1795 job: crate::message::OfflineBounceWork,
1796 ) {
1797 let track_name = job.track_name.clone();
1798 self.bounce_worker_tracks
1799 .insert(worker_index, track_name.clone());
1800 let worker = &self.workers[worker_index];
1801 if let Err(e) = worker.tx.send(Message::ProcessOfflineBounce(job)).await {
1802 self.bounce_worker_tracks.remove(&worker_index);
1803 self.offline_bounce_jobs.remove(&track_name);
1804 self.push_ready_worker(worker_index);
1805 self.notify_clients(Err(format!("Failed to schedule offline bounce: {e}")))
1806 .await;
1807 }
1808 }
1809
1810 async fn drain_pending_requests_if_idle(&mut self) {
1813 if self.offline_bounce_jobs.is_empty() {
1814 while let Some(next) = self.pending_requests.pop_front() {
1815 self.handle_request(next).await;
1816 }
1817 }
1818 }
1819
1820 pub(crate) async fn on_all_tracks_finished(&mut self) {
1821 let pending = std::mem::take(&mut self.pending_bounce_starts);
1824 for (worker_index, job) in pending {
1825 self.send_bounce_job(worker_index, job).await;
1826 }
1827 if self.transport_restart_pending {
1828 let state = self.state_snapshot.load_full();
1829 for track in state.tracks.values() {
1830 track.lock().take_hw_midi_out_events();
1831 }
1832 } else if self.hw_worker.is_some() {
1833 self.active_hw_notes_cycle_start = self.active_hw_notes_by_track.clone();
1834 let mut out_events = self.collect_hw_midi_output_events_by_device();
1835 if self.loop_enabled
1836 && let Some((_, loop_end)) = self.loop_range_samples
1837 {
1838 let cycle_end = self
1839 .transport_sample
1840 .saturating_add(self.current_cycle_samples());
1841 if self.transport_sample < loop_end && cycle_end >= loop_end {
1842 let wrap_frame = loop_end
1843 .saturating_sub(self.transport_sample)
1844 .min(self.current_cycle_samples())
1845 as u32;
1846 out_events.extend(self.note_off_events_for_active_snapshot(
1847 &self.active_hw_notes_cycle_start,
1848 wrap_frame,
1849 ));
1850 out_events.sort_by(|a, b| {
1851 a.event
1852 .frame
1853 .cmp(&b.event.frame)
1854 .then_with(|| a.device.cmp(&b.device))
1855 });
1856 }
1857 }
1858 self.pending_hw_midi_out_events_by_device.extend(out_events);
1859 } else {
1860 self.pending_hw_midi_out_events = self.collect_hw_midi_output_events();
1861 }
1862 self.request_hw_cycle().await;
1863 }
1864
1865 pub(crate) fn take_ready_worker_index(&mut self) -> Option<usize> {
1866 while !self.ready_workers.is_empty() {
1867 let worker_index = self.ready_workers.remove(0);
1868 if worker_index < self.workers.len() {
1869 return Some(worker_index);
1870 }
1871 }
1872 None
1873 }
1874
1875 pub(crate) fn push_ready_worker(&mut self, worker_index: usize) {
1876 self.ready_workers.push(worker_index);
1877 }
1878
1879 pub(crate) fn publish_track_meters(&mut self) {
1880 if !self.should_publish_track_meters() {
1881 return;
1882 }
1883 let tracks: Vec<(String, crate::state::TrackHandle)> = self
1884 .state_snapshot
1885 .load_full()
1886 .tracks
1887 .iter()
1888 .map(|(name, track)| (name.clone(), track.clone()))
1889 .collect();
1890 let mut snapshot = Vec::with_capacity(tracks.len());
1891 for (name, track) in &tracks {
1892 let linear = self
1893 .track_meter_linear_by_track
1894 .get(name)
1895 .cloned()
1896 .unwrap_or_else(|| track.lock().output_meter_linear());
1897 let output_db = linear
1898 .iter()
1899 .copied()
1900 .map(Self::meter_linear_to_db)
1901 .collect::<Vec<_>>();
1902 snapshot.push((name.clone(), output_db));
1903 }
1904 self.latest_track_meter_snapshot = Arc::new(snapshot);
1905 }
1906
1907 fn record_session_completed_clip_pass(
1908 &mut self,
1909 track_name: String,
1910 scene_index: usize,
1911 clip_id: String,
1912 pass_index: usize,
1913 start_sample: usize,
1914 length_samples: usize,
1915 ) {
1916 let key = (
1917 track_name.clone(),
1918 scene_index,
1919 clip_id.clone(),
1920 pass_index,
1921 start_sample,
1922 );
1923 if self.session_reported_clip_passes.insert(key) {
1924 self.session_completed_clip_passes
1925 .push(crate::meter::SessionCompletedClipPass {
1926 track_name,
1927 scene_index,
1928 clip_id,
1929 pass_index,
1930 start_sample,
1931 length_samples,
1932 });
1933 }
1934 }
1935
1936 fn record_completed_session_scene_span(
1937 &mut self,
1938 tracks: &HashMap<String, crate::state::TrackHandle>,
1939 scene_index: usize,
1940 previous_scene: Option<usize>,
1941 scene_start_sample: usize,
1942 elapsed_samples: usize,
1943 ) {
1944 for (track_name, track) in tracks {
1945 let track = track.lock();
1946 let slot = track.rt.session_slots.get(&scene_index);
1947 let prev_slot = previous_scene.and_then(|scene| track.rt.session_slots.get(&scene));
1948 let playing_clip_id = track
1949 .rt
1950 .playing_session_clips
1951 .last()
1952 .map(|clip| clip.clip_id.as_str());
1953 let clip_id = match (slot, prev_slot) {
1954 (Some(slot), _) if slot.play_enabled => Some(slot.clip_id.as_str()),
1955 (Some(slot), _) if slot.stop_enabled => None,
1956 (_, Some(prev_slot)) if prev_slot.play_enabled => Some(prev_slot.clip_id.as_str()),
1957 (_, Some(prev_slot)) if prev_slot.stop_enabled => None,
1958 _ => playing_clip_id,
1959 }
1960 .filter(|clip_id| !clip_id.is_empty());
1961 let Some(clip_id) = clip_id else { continue };
1962 let clip_length = track
1963 .session_clip_length(clip_id, crate::kind::Kind::Audio)
1964 .or_else(|| track.session_clip_length(clip_id, crate::kind::Kind::MIDI))
1965 .unwrap_or(0);
1966 if clip_length == 0 {
1967 continue;
1968 }
1969 let completed_passes = elapsed_samples / clip_length;
1970 for pass_index in 0..completed_passes {
1971 self.record_session_completed_clip_pass(
1972 track_name.clone(),
1973 scene_index,
1974 clip_id.to_string(),
1975 pass_index,
1976 scene_start_sample.saturating_add(pass_index * clip_length),
1977 clip_length,
1978 );
1979 }
1980 }
1981 }
1982
1983 pub(crate) async fn publish_session_runtime_reports(&mut self) {
1984 let mut current = HashMap::<(String, usize), (SessionSlotState, usize, usize)>::new();
1985 {
1986 let state = self.state_snapshot.load_full();
1987 if let Some((queued_scene, queued_launch_at)) = self.session_scene_queue {
1988 let still_pending = state.tracks.values().any(|track| {
1996 let track = track.lock();
1997 track.rt.pending_session_launches.iter().any(|launch| {
1998 launch.scene_index == queued_scene
1999 && launch.launch_at_sample == queued_launch_at
2000 }) || track
2001 .rt
2002 .playing_session_clips
2003 .iter()
2004 .any(|clip| clip.stop_at_sample == Some(queued_launch_at))
2005 });
2006 if !still_pending && self.session_transport_sample >= queued_launch_at {
2007 if let Some(current_scene) = self.session_current_scene {
2008 let scene_start = self.session_current_scene_start_sample;
2009 let elapsed = queued_launch_at.saturating_sub(scene_start);
2010 self.record_completed_session_scene_span(
2011 &state.tracks,
2012 current_scene,
2013 self.session_current_scene_previous_scene,
2014 scene_start,
2015 elapsed,
2016 );
2017 }
2018 self.session_scene_queue = None;
2019 self.session_current_scene_previous_scene = self.session_current_scene;
2020 self.session_current_scene = Some(queued_scene);
2021 self.session_current_scene_start_sample = queued_launch_at;
2022 self.session_current_scene_length_samples =
2023 self.session_scene_queue_length_samples;
2024 self.session_scene_queue_length_samples = 0;
2025 }
2026 }
2027 if let Some(current_scene) = self.session_current_scene
2028 && self.session_current_scene_length_samples > 0
2029 {
2030 self.record_completed_session_scene_span(
2031 &state.tracks,
2032 current_scene,
2033 self.session_current_scene_previous_scene,
2034 self.session_current_scene_start_sample,
2035 self.session_transport_sample
2036 .saturating_sub(self.session_current_scene_start_sample),
2037 );
2038 }
2039 for (track_name, track) in &state.tracks {
2040 let track = track.lock();
2041 for launch in &track.rt.pending_session_launches {
2042 current.insert(
2043 (track_name.clone(), launch.scene_index),
2044 (SessionSlotState::Queued, 0, 0),
2045 );
2046 }
2047 for clip in &track.rt.playing_session_clips {
2048 if self.session_current_scene_length_samples == 0
2049 && let Some(clip_length) =
2050 track.session_clip_length(&clip.clip_id, clip.kind)
2051 {
2052 let scene_report = self
2053 .session_current_scene
2054 .map(|_| self.session_current_scene_length_samples)
2055 .filter(|length| *length > 0)
2056 .map(|length| {
2057 (
2058 self.session_current_scene.unwrap_or(clip.scene_index),
2059 self.session_current_scene_start_sample,
2060 self.session_transport_sample
2061 .saturating_sub(self.session_current_scene_start_sample),
2062 length,
2063 )
2064 });
2065 let launch_sample = self
2066 .session_transport_sample
2067 .saturating_sub(clip.elapsed_samples);
2068 let (
2069 report_scene_index,
2070 report_start_sample,
2071 report_elapsed,
2072 report_length,
2073 ) = scene_report.unwrap_or((
2074 clip.scene_index,
2075 launch_sample,
2076 clip.elapsed_samples,
2077 clip_length,
2078 ));
2079 if report_length == 0 {
2080 continue;
2081 }
2082 let completed_passes = report_elapsed / report_length;
2083 for pass_index in 0..completed_passes {
2084 self.record_session_completed_clip_pass(
2085 track_name.clone(),
2086 report_scene_index,
2087 clip.clip_id.clone(),
2088 pass_index,
2089 report_start_sample.saturating_add(pass_index * report_length),
2090 report_length,
2091 );
2092 }
2093 }
2094 current.insert(
2095 (track_name.clone(), clip.scene_index),
2096 (
2097 SessionSlotState::Playing,
2098 clip.play_position_samples,
2099 clip.elapsed_samples,
2100 ),
2101 );
2102 }
2103 }
2104 }
2105
2106 if self
2107 .last_session_report_publish
2108 .is_some_and(|t| t.elapsed() < Self::SESSION_RUNTIME_REPORT_INTERVAL)
2109 {
2110 return;
2111 }
2112
2113 let snapshot = self.session_runtime_snapshot_producer.write_buffer();
2114 snapshot.session_sample = self.session_transport_sample;
2115 snapshot.slots.clear();
2116 snapshot.slots.extend(current.iter().map(
2117 |((track_name, scene_index), (state, play_position_samples, elapsed_samples))| {
2118 crate::meter::SessionRuntimeSlotSnapshot {
2119 track_name: track_name.clone(),
2120 scene_index: *scene_index,
2121 state: *state,
2122 play_position_samples: *play_position_samples,
2123 elapsed_samples: *elapsed_samples,
2124 }
2125 },
2126 ));
2127 snapshot.completed_clip_passes = self.session_completed_clip_passes.clone();
2128 snapshot.current_scene = self.session_current_scene;
2129 self.session_runtime_snapshot_producer.publish();
2130 self.last_session_report_publish = Some(Instant::now());
2131 }
2132
2133 pub(crate) async fn publish_clap_state_dirty(&mut self) {
2134 let tracks: Vec<(String, crate::state::TrackHandle)> = self
2135 .state_snapshot
2136 .load_full()
2137 .tracks
2138 .iter()
2139 .map(|(name, track)| (name.clone(), track.clone()))
2140 .collect();
2141 for (track_name, track) in &tracks {
2142 let dirty = track.lock().take_dirty_clap_instances();
2143 for instance_id in dirty {
2144 self.notify_clients(Ok(Action::TrackClapStateDirty {
2145 track_name: track_name.clone(),
2146 instance_id,
2147 }))
2148 .await;
2149 }
2150 }
2151 }
2152
2153 pub(crate) fn reset_meters_after_stop(&mut self) {
2154 self.last_hw_out_meter_publish = None;
2155 self.last_track_meter_publish = None;
2156 self.last_meter_snapshot_publish = None;
2157 self.hw_out_peak_hold_linear.fill(0.0);
2158 #[cfg(any(target_os = "freebsd", target_os = "linux", target_os = "openbsd"))]
2159 {
2160 self.last_hw_out_meter_linear.clear();
2161 }
2162 let hw_channels = self.latest_hw_out_meter_db.len();
2163 self.latest_hw_out_meter_db = Arc::new(vec![-90.0; hw_channels]);
2164
2165 let tracks: Vec<(String, crate::state::TrackHandle)> = self
2166 .state_snapshot
2167 .load_full()
2168 .tracks
2169 .iter()
2170 .map(|(name, track)| (name.clone(), track.clone()))
2171 .collect();
2172 self.track_meter_linear_by_track.clear();
2173 let mut snapshot = Vec::with_capacity(tracks.len());
2174 for (name, track) in tracks {
2175 let mut t = track.lock();
2176 t.clear_output_meters();
2177 let width = t.output_meter_linear().len();
2178 let zero_linear = vec![0.0; width];
2179 self.track_meter_linear_by_track
2180 .insert(name.clone(), zero_linear);
2181 snapshot.push((name, vec![-90.0; width]));
2182 }
2183 self.latest_track_meter_snapshot = Arc::new(snapshot);
2184 self.publish_meter_snapshot();
2185 }
2186
2187 pub(crate) fn publish_meter_snapshot_if_due(&mut self) {
2188 let now = Instant::now();
2189 if self
2190 .last_meter_snapshot_publish
2191 .is_some_and(|last| now.duration_since(last) < Self::METER_PUBLISH_INTERVAL)
2192 {
2193 return;
2194 }
2195 self.last_meter_snapshot_publish = Some(now);
2196 self.publish_meter_snapshot();
2197 }
2198
2199 pub(crate) fn publish_meter_snapshot(&mut self) {
2200 self.update_hw_out_lufs_readout();
2203
2204 let snapshot = self.meter_snapshot_producer.write_buffer();
2205 snapshot.hw_out_db.clear();
2206 snapshot
2207 .hw_out_db
2208 .extend(self.latest_hw_out_meter_db.iter().copied());
2209 snapshot.hw_out_lufs = self.latest_hw_out_lufs;
2210 snapshot.track_meters.clear();
2211 snapshot
2212 .track_meters
2213 .extend(self.latest_track_meter_snapshot.iter().cloned());
2214 self.meter_snapshot_producer.publish();
2215 }
2216
2217 pub(crate) fn publish_transport_snapshot(&mut self) {
2218 let snapshot = self.transport_snapshot_producer.write_buffer();
2219 snapshot.sample = self.transport_sample;
2220 snapshot.tempo_bpm = self.tempo_bpm;
2221 snapshot.playing = self.playing;
2222 snapshot.transport_running = self.transport_running;
2223 snapshot.tsig_num = self.tsig_num;
2224 snapshot.tsig_denom = self.tsig_denom;
2225 self.transport_snapshot_producer.publish();
2226 }
2227
2228 pub(crate) async fn handle_request(&mut self, a: Action) {
2229 match a {
2230 Action::Log { source, message } => {
2231 self.notify_clients(Ok(Action::Log { source, message }))
2232 .await;
2233 }
2234 Action::Undo => {
2235 let actions = match self.history.undo() {
2236 Some(actions) => actions,
2237 None => {
2238 self.notify_clients(Ok(Action::Undo)).await;
2239 self.notify_clients(Ok(Action::HistoryState {
2240 dirty: self.history.is_dirty(),
2241 }))
2242 .await;
2243 return;
2244 }
2245 };
2246
2247 let was_suspended = self.history_suspended;
2248 self.history_suspended = true;
2249 for action in actions {
2250 self.handle_request_inner(action, false).await;
2251 }
2252 self.history_suspended = was_suspended;
2253 self.notify_clients(Ok(Action::Undo)).await;
2254 self.notify_clients(Ok(Action::HistoryState {
2255 dirty: self.history.is_dirty(),
2256 }))
2257 .await;
2258 }
2259 Action::Redo => {
2260 let actions = match self.history.redo() {
2261 Some(actions) => actions,
2262 None => {
2263 self.notify_clients(Ok(Action::Redo)).await;
2264 self.notify_clients(Ok(Action::HistoryState {
2265 dirty: self.history.is_dirty(),
2266 }))
2267 .await;
2268 return;
2269 }
2270 };
2271
2272 let was_suspended = self.history_suspended;
2273 self.history_suspended = true;
2274 for action in actions {
2275 self.handle_request_inner(action, false).await;
2276 }
2277 self.history_suspended = was_suspended;
2278 self.notify_clients(Ok(Action::Redo)).await;
2279 self.notify_clients(Ok(Action::HistoryState {
2280 dirty: self.history.is_dirty(),
2281 }))
2282 .await;
2283 }
2284 Action::ApplyGroupedActions(actions) => {
2285 self.handle_request_inner(Action::BeginHistoryGroup, true)
2286 .await;
2287 for action in actions {
2288 self.handle_request_inner(action, true).await;
2289 }
2290 self.handle_request_inner(Action::EndHistoryGroup, true)
2291 .await;
2292 }
2293 Action::Session(_) => {
2294 self.handle_request_inner(a, false).await;
2295 }
2296 other => {
2297 self.handle_request_inner(other, true).await;
2298 }
2299 }
2300 self.publish_state_snapshot();
2301 }
2302
2303 pub(crate) async fn handle_quit(&mut self, a: Action) {
2304 self.flush_recordings().await;
2305 if let Some(mut worker) = self.hw_worker.take() {
2313 let panic_events = self.panic_events_for_all_hw_midi_outputs();
2316 if !panic_events.is_empty() {
2317 let _ = worker.tx.send(Message::HWMidiOutEvents(panic_events)).await;
2318 }
2319 if let Err(e) = worker.tx.send(Message::Request(a.clone())).await {
2322 error!("Error sending quit message to HW worker: {e}");
2323 }
2324 if let Some(handle) = worker.handle.take() {
2325 handle
2326 .await
2327 .unwrap_or_else(|e| error!("Error waiting for HW worker to quit: {e}"));
2328 }
2329 }
2330 if let Some(hw) = self.hw_driver.as_mut() {
2336 hw.close_fds();
2337 }
2338 if let Some(midi_hub) = self.midi_hub.as_mut() {
2339 midi_hub.close_all();
2340 }
2341 self.hw_driver = None;
2342 self.hw_driver_info = None;
2343 self.hw_input_ports.clear();
2344 self.hw_output_ports.clear();
2345 self.notify_clients(Ok(Action::Quit)).await;
2346 self.ready_workers.clear();
2347 while !self.workers.is_empty() {
2348 let mut worker = self.workers.remove(0);
2349 if let Err(e) = worker.tx.send(Message::Request(a.clone())).await {
2350 error!("Error sending quit message to worker: {e}");
2351 }
2352 if let Some(handle) = worker.handle.take() {
2353 handle
2354 .await
2355 .unwrap_or_else(|e| error!("Error waiting for worker to quit: {e}"));
2356 }
2357 }
2358 #[cfg(unix)]
2359 {
2360 self.jack_runtime = None;
2361 }
2362 self.osc_server = None;
2363 }
2364
2365 #[inline]
2366 pub(crate) fn box_bool<'a>(
2367 fut: impl std::future::Future<Output = bool> + Send + 'a,
2368 ) -> std::pin::Pin<Box<dyn std::future::Future<Output = bool> + Send + 'a>> {
2369 Box::pin(fut)
2370 }
2371
2372 pub(crate) async fn handle_request_inner(
2373 &mut self,
2374 mut action_to_process: Action,
2375 record_history: bool,
2376 ) {
2377 let a = action_to_process.clone();
2378 let suppress_timing_history = self.playing
2379 && matches!(
2380 &action_to_process,
2381 Action::SetTempo(_) | Action::SetTimeSignature { .. } | Action::SetTempoMap { .. }
2382 );
2383 let mut inverse_actions = self.prepare_inverse_actions(
2384 &action_to_process,
2385 record_history,
2386 suppress_timing_history,
2387 );
2388
2389 match action_to_process {
2390 Action::Play => {
2391 if Self::box_bool(self.handle_play(a.clone())).await {
2392 return;
2393 }
2394 }
2395 Action::Pause => {
2396 if Self::box_bool(self.handle_pause(a.clone())).await {
2397 return;
2398 }
2399 }
2400 Action::Stop => {
2401 if Self::box_bool(self.handle_stop(a.clone())).await {
2402 return;
2403 }
2404 }
2405 Action::SessionPlay => {
2406 if Self::box_bool(self.handle_session_play(a.clone())).await {
2407 return;
2408 }
2409 }
2410 Action::JumpToEnd => {
2411 self.transport_sample = self.normalize_transport_sample(self.session_end_sample());
2412 self.publish_transport_snapshot();
2413 self.notify_clients(Ok(Action::TransportPosition(self.transport_sample)))
2414 .await;
2415 }
2416 Action::Panic => {
2417 if Self::box_bool(self.handle_panic(a.clone())).await {
2418 return;
2419 }
2420 }
2421 Action::Session(ref session_action) => {
2422 self.handle_session_action(session_action.clone()).await;
2423 }
2424 Action::SessionRuntimeReport { .. } => {}
2425 Action::SessionMidiLearnTriggered { .. } => {}
2426 Action::SetClipPlaybackEnabled(enabled) => {
2427 self.clip_playback_enabled = enabled;
2428 for track in self.state_snapshot.load_full().tracks.values() {
2429 track.lock().set_clip_playback_enabled(enabled);
2430 }
2431 }
2432 Action::SetSessionClipPlaybackEnabled(enabled) => {
2433 self.session_clip_playback_enabled = enabled;
2434 for track in self.state_snapshot.load_full().tracks.values() {
2435 track.lock().set_session_clip_playback_enabled(enabled);
2436 }
2437 }
2438 Action::TransportPosition(..) => {
2439 if Self::box_bool(self.handle_transport_position(a.clone())).await {
2440 return;
2441 }
2442 }
2443 Action::SetLoopEnabled(enabled) => {
2444 self.loop_enabled = enabled && self.loop_range_samples.is_some();
2445 self.notified_loop_wrap_sample = None;
2446 }
2447 Action::SetLoopRange(..) => {
2448 if Self::box_bool(self.handle_set_loop_range(a.clone())).await {
2449 return;
2450 }
2451 }
2452 Action::SetPunchEnabled(enabled) => {
2453 self.punch_enabled = enabled && self.punch_range_samples.is_some();
2454 }
2455 Action::SetPunchRange(range) => {
2456 self.punch_range_samples = range.and_then(|(start, end)| {
2457 if end > start {
2458 Some((start, end))
2459 } else {
2460 None
2461 }
2462 });
2463 self.punch_enabled = self.punch_range_samples.is_some();
2464 }
2465 Action::SetMetronomeEnabled(enabled) => {
2466 self.metronome_enabled = enabled;
2467 if enabled {
2468 self.ensure_metronome_track().await;
2469 }
2470 if let Some(track) = self
2471 .state_snapshot
2472 .load_full()
2473 .tracks
2474 .get(Self::METRONOME_TRACK)
2475 .cloned()
2476 {
2477 track.lock().set_metronome_enabled(enabled);
2478 }
2479 }
2480 Action::SetTempo(bpm) => {
2481 self.tempo_bpm = bpm.max(1.0);
2482 self.publish_transport_snapshot();
2483 }
2484 Action::SetTimeSignature {
2485 numerator,
2486 denominator,
2487 } => {
2488 self.tsig_num = numerator.max(1);
2489 self.tsig_denom = denominator.max(1);
2490 self.publish_transport_snapshot();
2491 }
2492 Action::SetTempoMap {
2493 ref tempo_points,
2494 ref time_signature_points,
2495 } => {
2496 self.tempo_points = tempo_points.clone();
2497 self.time_signature_points = time_signature_points.clone();
2498 self.update_global_tempo_from_map();
2499 self.publish_transport_snapshot();
2500 }
2501 Action::SetOscEnabled(enabled) => {
2502 if let Err(err) = self.set_osc_enabled_with(enabled, OscServer::start) {
2503 self.notify_clients(Err(err)).await;
2504 }
2505 }
2506 Action::SetRecordEnabled(..) => {
2507 if Self::box_bool(self.handle_set_record_enabled(a.clone())).await {
2508 return;
2509 }
2510 }
2511 Action::SetModulators(ref modulators) => {
2512 self.modulators = modulators.clone();
2513 let echoes = self.apply_modulators(self.active_transport_sample());
2514 for action in echoes {
2515 self.notify_clients(Ok(action)).await;
2516 }
2517 }
2518 Action::SetTrackAutomationLanes {
2519 ref track_name,
2520 ref lanes,
2521 mode,
2522 } => {
2523 if let Some(track) = self.state_snapshot.load_full().tracks.get(track_name) {
2524 let mut track = track.lock();
2525 track.automation_lanes = lanes.clone();
2526 track.set_automation_mode(mode);
2527 }
2528 }
2529 Action::TrackAutomationToggleLane { .. } => {
2530 if Self::box_bool(self.handle_track_automation_toggle_lane(a.clone())).await {
2531 return;
2532 }
2533 }
2534 Action::TrackAutomationInsertPoint { .. } => {
2535 if Self::box_bool(self.handle_track_automation_insert_point(a.clone())).await {
2536 return;
2537 }
2538 }
2539 Action::TrackAutomationDeletePoint { .. } => {
2540 if Self::box_bool(self.handle_track_automation_delete_point(a.clone())).await {
2541 return;
2542 }
2543 }
2544 Action::TrackAutomationSetMode {
2545 ref track_name,
2546 mode,
2547 } => {
2548 if let Some(track) = self
2549 .state_snapshot
2550 .load_full()
2551 .tracks
2552 .get(track_name)
2553 .cloned()
2554 {
2555 track.lock().set_automation_mode(mode);
2556 }
2557 }
2558 Action::RequestTrackList => {
2559 let names: Vec<String> = self
2560 .state_snapshot
2561 .load_full()
2562 .tracks
2563 .keys()
2564 .cloned()
2565 .collect();
2566 self.notify_clients(Ok(Action::TrackList(names))).await;
2567 }
2568 Action::TrackList(_) => {}
2569 Action::RequestTransportState => {
2570 self.notify_clients(Ok(Action::TransportState {
2571 sample: self.transport_sample,
2572 tempo_bpm: self.tempo_bpm,
2573 playing: self.playing,
2574 paused: !self.transport_running && self.playing,
2575 tsig_num: self.tsig_num,
2576 tsig_denom: self.tsig_denom,
2577 }))
2578 .await;
2579 }
2580 Action::TransportState { .. } => {}
2581 Action::SetStepRecording(enabled) => {
2582 self.step_recording_enabled = enabled;
2583 }
2584 Action::BeginHistoryGroup if self.history_group.is_none() => {
2585 self.history_group = Some(UndoEntry {
2586 forward_actions: vec![],
2587 inverse_actions: vec![],
2588 });
2589 }
2590 Action::EndHistoryGroup => {
2591 if let Some(mut group) = self.history_group.take()
2592 && !group.forward_actions.is_empty()
2593 && !group.inverse_actions.is_empty()
2594 {
2595 let mut add_tracks = Vec::new();
2596 let mut connections = Vec::new();
2597 let mut rest = Vec::new();
2598 for action in group.inverse_actions {
2599 if matches!(action, Action::AddTrack { .. }) {
2600 add_tracks.push(action);
2601 } else if matches!(action, Action::Connect { .. }) {
2602 connections.push(action);
2603 } else {
2604 rest.push(action);
2605 }
2606 }
2607 group.inverse_actions = add_tracks;
2608 group.inverse_actions.extend(rest);
2609 group.inverse_actions.extend(connections);
2610 self.history.record(group);
2611 }
2612 }
2613 Action::SetSessionPath(ref path) => {
2614 self.session_dir = Some(Path::new(path).to_path_buf());
2615 self.ensure_session_subdirs();
2616 #[cfg(all(unix, not(target_os = "macos")))]
2617 let _lv2_dir = self.session_plugins_dir();
2618 for track in self.state_snapshot.load_full().tracks.values() {
2619 track.lock().set_session_base_dir(self.session_dir.clone());
2620 }
2621 }
2622 Action::MarkHistorySavePoint => {
2623 self.history.mark_save_point();
2624 self.notify_clients(Ok(Action::HistoryState {
2625 dirty: self.history.is_dirty(),
2626 }))
2627 .await;
2628 }
2629 Action::ClearHistory => {
2630 self.history.clear();
2631 self.history.mark_save_point();
2632 }
2633 Action::BeginSessionRestore => {
2634 self.history_suspended = true;
2635 self.history.clear();
2636 }
2637 Action::EndSessionRestore => {
2638 self.history.clear();
2639 self.history_suspended = false;
2640 self.preload_track_clips_spawn();
2641 }
2642 Action::Quit => {
2643 self.handle_quit(a.clone()).await;
2644 return;
2645 }
2646 Action::AddTrack {
2647 ref name,
2648 audio_ins,
2649 midi_ins,
2650 audio_outs,
2651 midi_outs,
2652 folder,
2653 } => {
2654 self.handle_add_track(
2655 name.clone(),
2656 audio_ins,
2657 midi_ins,
2658 audio_outs,
2659 midi_outs,
2660 folder,
2661 )
2662 .await;
2663 }
2664 Action::TrackAddAudioInput(..) => {
2665 if Self::box_bool(self.handle_track_add_audio_input(a.clone())).await {
2666 return;
2667 }
2668 }
2669 Action::TrackAddAudioOutput(..) => {
2670 if Self::box_bool(self.handle_track_add_audio_output(a.clone())).await {
2671 return;
2672 }
2673 }
2674 Action::TrackRemoveAudioInput(..) => {
2675 if Self::box_bool(self.handle_track_remove_audio_input(a.clone())).await {
2676 return;
2677 }
2678 }
2679 Action::TrackRemoveAudioOutput(..) => {
2680 if Self::box_bool(self.handle_track_remove_audio_output(a.clone())).await {
2681 return;
2682 }
2683 }
2684 Action::RenameTrack { .. } => {
2685 if Self::box_bool(self.handle_rename_track(a.clone())).await {
2686 return;
2687 }
2688 }
2689 Action::RemoveTrack(ref name) => {
2690 self.handle_remove_track(name.clone(), record_history).await;
2691 inverse_actions = None;
2692 }
2693 Action::TrackLevel(ref name, level) => {
2694 if name == "hw:out" {
2695 self.hw_out_level_db = level;
2696 } else if let Some(track) = self.state_snapshot.load_full().tracks.get(name) {
2697 track.lock().set_level(level);
2698 }
2699 }
2700 Action::TrackBalance(ref name, balance) => {
2701 if name == "hw:out" {
2702 self.hw_out_balance = balance.clamp(-1.0, 1.0);
2703 } else if let Some(track) = self.state_snapshot.load_full().tracks.get(name) {
2704 track.lock().set_balance(balance);
2705 }
2706 }
2707 Action::TrackAutomationLevel(ref name, level) => {
2708 tracing::debug!(%name, level, "engine received TrackAutomationLevel");
2709 if name == "hw:out" {
2710 self.hw_out_level_db = level;
2711 } else if let Some(track) = self.state_snapshot.load_full().tracks.get(name) {
2712 track.lock().set_level(level);
2713 }
2714 }
2715 Action::TrackAutomationBalance(ref name, balance) => {
2716 if name == "hw:out" {
2717 self.hw_out_balance = balance.clamp(-1.0, 1.0);
2718 } else if let Some(track) = self.state_snapshot.load_full().tracks.get(name) {
2719 track.lock().set_balance(balance);
2720 }
2721 }
2722 Action::TrackMidiCc { .. } => {
2723 if Self::box_bool(self.handle_track_midi_cc(a.clone())).await {
2724 return;
2725 }
2726 }
2727 Action::RequestMeterSnapshot => {
2728 self.notify_clients(Ok(Action::MeterSnapshot {
2729 hw_out_db: self.latest_hw_out_meter_db.clone(),
2730 hw_out_lufs: self.latest_hw_out_lufs,
2731 track_meters: self.latest_track_meter_snapshot.clone(),
2732 }))
2733 .await;
2734 return;
2735 }
2736 Action::TrackMeters { .. } => {}
2737 Action::MeterSnapshot { .. } => {}
2738 Action::TrackToggleArm(..) => {
2739 if Self::box_bool(self.handle_track_toggle_arm(a.clone())).await {
2740 return;
2741 }
2742 }
2743 Action::TrackToggleMute(ref name) => {
2744 if name == "hw:out" {
2745 self.hw_out_muted = !self.hw_out_muted;
2746 } else if let Some(track) = self.state_snapshot.load_full().tracks.get(name) {
2747 track.lock().mute();
2748 }
2749 }
2750 Action::TrackTogglePhase(ref name) => {
2751 if let Some(track) = self.state_snapshot.load_full().tracks.get(name) {
2752 track.lock().invert_phase();
2753 }
2754 }
2755 Action::TrackToggleSolo(ref name) => {
2756 if name == "hw:out" {
2757 return;
2758 }
2759 if let Some(track) = self.state_snapshot.load_full().tracks.get(name) {
2760 track.lock().solo();
2761 }
2762 }
2763 Action::TrackToggleMaster(ref name) => {
2764 if let Some(track) = self.state_snapshot.load_full().tracks.get(name) {
2765 track.lock().toggle_master();
2766 }
2767 }
2768 Action::TrackToggleInputMonitor {
2769 ref track_name,
2770 lane,
2771 } => {
2772 if let Some(track) = self.state_snapshot.load_full().tracks.get(track_name) {
2773 track.lock().toggle_input_monitor(lane);
2774 }
2775 }
2776 Action::TrackToggleDiskMonitor {
2777 ref track_name,
2778 lane,
2779 } => {
2780 if let Some(track) = self.state_snapshot.load_full().tracks.get(track_name) {
2781 track.lock().toggle_disk_monitor(lane);
2782 }
2783 }
2784 Action::TrackToggleMidiInputMonitor {
2785 ref track_name,
2786 lane,
2787 } => {
2788 if let Some(track) = self.state_snapshot.load_full().tracks.get(track_name) {
2789 track.lock().toggle_midi_input_monitor(lane);
2790 }
2791 }
2792 Action::TrackToggleMidiDiskMonitor {
2793 ref track_name,
2794 lane,
2795 } => {
2796 if let Some(track) = self.state_snapshot.load_full().tracks.get(track_name) {
2797 track.lock().toggle_midi_disk_monitor(lane);
2798 }
2799 }
2800 Action::TrackSetColor {
2801 ref track_name,
2802 color,
2803 } => {
2804 if let Some(track) = self.state_snapshot.load_full().tracks.get(track_name) {
2805 track.lock().color = color;
2806 }
2807 }
2808 Action::TrackArmMidiLearn {
2809 ref track_name,
2810 target,
2811 } => {
2812 if let Err(e) = self.track_handle_or_err(track_name) {
2813 self.notify_clients(Err(e)).await;
2814 return;
2815 }
2816 self.pending_midi_learn = Some((track_name.clone(), target, None));
2817 }
2818 Action::GlobalArmMidiLearn { target } => {
2819 self.pending_global_midi_learn = Some(target);
2820 }
2821 Action::SessionArmMidiLearn { ref target } => {
2822 self.pending_session_midi_learn = Some(target.clone());
2823 }
2824 Action::TrackSetMidiLearnBinding { .. } => {
2825 if Self::box_bool(self.handle_track_set_midi_learn_binding(a.clone())).await {
2826 return;
2827 }
2828 }
2829 Action::SetGlobalMidiLearnBinding { .. } => {
2830 if Self::box_bool(self.handle_set_global_midi_learn_binding(a.clone())).await {
2831 return;
2832 }
2833 }
2834 Action::SetSessionMidiLearnBinding { .. } => {
2835 if Self::box_bool(self.handle_set_session_midi_learn_binding(a.clone())).await {
2836 return;
2837 }
2838 }
2839 Action::TrackSetFolder { .. } => {
2840 if Self::box_bool(self.handle_track_set_folder(a.clone())).await {
2841 return;
2842 }
2843 }
2844 Action::TrackSetParent {
2845 ref track_name,
2846 ref parent_name,
2847 } => {
2848 self.handle_track_set_parent(track_name.as_str(), parent_name.as_deref())
2849 .await;
2850 }
2851 Action::TrackToggleFolder { .. } => {
2852 if Self::box_bool(self.handle_track_toggle_folder(a.clone())).await {
2853 return;
2854 }
2855 }
2856 Action::TrackSetMidiLaneChannel { .. } => {
2857 if Self::box_bool(self.handle_track_set_midi_lane_channel(a.clone())).await {
2858 return;
2859 }
2860 }
2861 Action::TrackSetFrozen { .. } => {
2862 if Self::box_bool(self.handle_track_set_frozen(a.clone())).await {
2863 return;
2864 }
2865 }
2866 Action::TrackSetSessionSlot { .. } => {
2867 if Self::box_bool(self.handle_track_set_session_slot(a.clone())).await {
2868 return;
2869 }
2870 }
2871 Action::TrackSetSessionSlotPlayEnabled { .. } => {
2872 if self
2873 .handle_track_set_session_slot_play_enabled(a.clone())
2874 .await
2875 {
2876 return;
2877 }
2878 }
2879 Action::TrackSetSessionSlotStopEnabled { .. } => {
2880 if self
2881 .handle_track_set_session_slot_stop_enabled(a.clone())
2882 .await
2883 {
2884 return;
2885 }
2886 }
2887 Action::TrackOfflineBounce { .. } => {
2888 self.handle_track_offline_bounce(action_to_process).await;
2889 return;
2890 }
2891 Action::TrackOfflineBounceCancel { .. } => {}
2892 Action::TrackOfflineBounceCancelAll => {}
2893 Action::TrackOfflineBounceCanceled { .. } => {}
2894 Action::TrackOfflineBounceProgress { .. } => {}
2895 Action::PianoKey {
2896 ref track_name,
2897 note,
2898 velocity,
2899 on,
2900 } => {
2901 if let Some(track) = self.state_snapshot.load_full().tracks.get(track_name) {
2902 let status = if on { 0x90 } else { 0x80 };
2903 let event = MidiEvent::new(0, vec![status, note.min(127), velocity.min(127)]);
2904 track.lock().push_hw_midi_events(&[event]);
2905 }
2906 }
2907 Action::ModifyMidiNotes { .. }
2908 | Action::ModifyMidiControllers { .. }
2909 | Action::DeleteMidiControllers { .. }
2910 | Action::InsertMidiControllers { .. }
2911 | Action::DeleteMidiNotes { .. }
2912 | Action::InsertMidiNotes { .. } => {
2913 if let Err(e) = self.apply_midi_edit_action(&action_to_process) {
2914 self.notify_clients(Err(e)).await;
2915 return;
2916 }
2917 }
2918 Action::SetMidiSysExEvents { .. } => {
2919 if let Err(e) = self.apply_midi_edit_action(&action_to_process) {
2920 self.notify_clients(Err(e)).await;
2921 return;
2922 }
2923 }
2924 Action::TrackClearDefaultPassthrough { .. } => {
2925 if Self::box_bool(self.handle_track_clear_default_passthrough(a.clone())).await {
2926 return;
2927 }
2928 }
2929 Action::TrackClearPlugins { .. } => {
2930 if Self::box_bool(self.handle_track_clear_plugins(a.clone())).await {
2931 return;
2932 }
2933 }
2934 #[cfg(all(unix, not(target_os = "macos")))]
2935 Action::TrackSetLv2PluginState { .. } => {
2936 if Self::box_bool(self.handle_track_set_lv2_plugin_state(a.clone())).await {
2937 return;
2938 }
2939 }
2940 #[cfg(all(unix, not(target_os = "macos")))]
2941 Action::ClipSetLv2PluginState { ref track_name, .. } => {
2942 self.notify_clients(Err(format!(
2943 "Track '{}': clip LV2 plugin state changes are not supported",
2944 track_name
2945 )))
2946 .await;
2947 }
2948 Action::TrackGetClapNoteNames { .. } => {
2949 if Self::box_bool(self.handle_track_get_clap_note_names(a.clone())).await {
2950 return;
2951 }
2952 }
2953 #[cfg(all(unix, not(target_os = "macos")))]
2954 Action::TrackGetLv2Midnam { .. } => {
2955 if Self::box_bool(self.handle_track_get_lv2_midnam(a.clone())).await {
2956 return;
2957 }
2958 }
2959 Action::TrackGetPluginGraph { .. } => {
2960 if Self::box_bool(self.handle_track_get_plugin_graph(a.clone())).await {
2961 return;
2962 }
2963 }
2964 Action::TrackPluginGraph { .. } => {}
2965 Action::TrackConnectPluginAudio { .. } => {
2966 if Self::box_bool(self.handle_track_connect_plugin_audio(a.clone())).await {
2967 return;
2968 }
2969 }
2970 Action::TrackConnectPluginMidi { .. } => {
2971 if Self::box_bool(self.handle_track_connect_plugin_midi(a.clone())).await {
2972 return;
2973 }
2974 }
2975 Action::TrackDisconnectPluginAudio { .. } => {
2976 if Self::box_bool(self.handle_track_disconnect_plugin_audio(a.clone())).await {
2977 return;
2978 }
2979 }
2980 Action::TrackDisconnectPluginMidi { .. } => {
2981 if Self::box_bool(self.handle_track_disconnect_plugin_midi(a.clone())).await {
2982 return;
2983 }
2984 }
2985 Action::TrackConnectAudio { .. } => {
2986 if Self::box_bool(self.handle_track_connect_audio(a.clone())).await {
2987 return;
2988 }
2989 }
2990 Action::TrackDisconnectAudio { .. } => {
2991 if Self::box_bool(self.handle_track_disconnect_audio(a.clone())).await {
2992 return;
2993 }
2994 }
2995 Action::TrackConnectMidi { .. } => {
2996 if Self::box_bool(self.handle_track_connect_midi(a.clone())).await {
2997 return;
2998 }
2999 }
3000 Action::TrackDisconnectMidi { .. } => {
3001 if Self::box_bool(self.handle_track_disconnect_midi(a.clone())).await {
3002 return;
3003 }
3004 }
3005 #[cfg(all(unix, not(target_os = "macos")))]
3006 Action::ListLv2Plugins => {
3007 if Self::box_bool(self.handle_list_lv2_plugins(a.clone())).await {
3008 return;
3009 }
3010 }
3011 #[cfg(all(unix, not(target_os = "macos")))]
3012 Action::Lv2Plugins(_) => {}
3013 #[cfg(all(unix, not(target_os = "macos")))]
3014 Action::Lv2PluginsUnavailable { .. } => {}
3015 Action::ListVst3Plugins => {
3016 if Self::box_bool(self.handle_list_vst3_plugins(a.clone())).await {
3017 return;
3018 }
3019 }
3020 Action::Vst3Plugins(_) => {}
3021 Action::Vst3PluginsUnavailable { .. } => {}
3022 Action::ListClapPlugins => {
3023 if Self::box_bool(self.handle_list_clap_plugins(a.clone())).await {
3024 return;
3025 }
3026 }
3027 Action::ListClapPluginsWithCapabilities => {
3028 if self
3029 .handle_list_clap_plugins_with_capabilities(a.clone())
3030 .await
3031 {
3032 return;
3033 }
3034 }
3035 Action::ClapPlugins(_) => {}
3036 Action::ClapPluginsUnavailable { .. } => {}
3037 Action::TrackLoadClapPlugin {
3038 ref track_name,
3039 ref plugin_id,
3040 instance_id,
3041 } => {
3042 if self
3043 .handle_track_load_clap_plugin(
3044 track_name.as_str(),
3045 plugin_id.as_str(),
3046 instance_id,
3047 )
3048 .await
3049 {
3050 return;
3051 }
3052 }
3053 Action::TrackUnloadClapPlugin {
3054 ref track_name,
3055 ref plugin_id,
3056 } => {
3057 if self
3058 .handle_track_unload_clap_plugin(track_name.as_str(), plugin_id.as_str())
3059 .await
3060 {
3061 return;
3062 }
3063 }
3064 Action::TrackUnloadClapPluginInstance {
3065 ref track_name,
3066 instance_id,
3067 } => {
3068 if self
3069 .handle_track_unload_clap_plugin_instance(track_name.as_str(), instance_id)
3070 .await
3071 {
3072 return;
3073 }
3074 }
3075 Action::TrackShowClapGui { .. } => {
3076 if Self::box_bool(self.handle_track_show_clap_gui(a.clone())).await {
3077 return;
3078 }
3079 }
3080 Action::TrackLoadVst3Plugin {
3081 ref track_name,
3082 ref plugin_id,
3083 instance_id,
3084 } => {
3085 if self
3086 .handle_track_load_vst3_plugin(
3087 track_name.as_str(),
3088 plugin_id.as_str(),
3089 instance_id,
3090 )
3091 .await
3092 {
3093 return;
3094 }
3095 }
3096 Action::TrackUnloadVst3Plugin {
3097 ref track_name,
3098 ref plugin_id,
3099 } => {
3100 if self
3101 .handle_track_unload_vst3_plugin(track_name.as_str(), plugin_id.as_str())
3102 .await
3103 {
3104 return;
3105 }
3106 }
3107 Action::TrackUnloadVst3PluginInstance {
3108 ref track_name,
3109 instance_id,
3110 } => {
3111 if self
3112 .handle_track_unload_vst3_plugin_instance(track_name.as_str(), instance_id)
3113 .await
3114 {
3115 return;
3116 }
3117 }
3118 Action::TrackShowVst3Gui { .. } => {
3119 if Self::box_bool(self.handle_track_show_vst3_gui(a.clone())).await {
3120 return;
3121 }
3122 }
3123 #[cfg(all(unix, not(target_os = "macos")))]
3124 Action::TrackLoadLv2Plugin {
3125 ref track_name,
3126 ref plugin_uri,
3127 instance_id,
3128 } => {
3129 if self
3130 .handle_track_load_lv2_plugin(
3131 track_name.as_str(),
3132 plugin_uri.as_str(),
3133 instance_id,
3134 )
3135 .await
3136 {
3137 return;
3138 }
3139 }
3140 #[cfg(all(unix, not(target_os = "macos")))]
3141 Action::TrackUnloadLv2Plugin {
3142 ref track_name,
3143 ref plugin_uri,
3144 } => {
3145 if self
3146 .handle_track_unload_lv2_plugin(track_name.as_str(), plugin_uri.as_str())
3147 .await
3148 {
3149 return;
3150 }
3151 }
3152 #[cfg(all(unix, not(target_os = "macos")))]
3153 Action::TrackUnloadLv2PluginInstance {
3154 ref track_name,
3155 instance_id,
3156 } => {
3157 if self
3158 .handle_track_unload_lv2_plugin_instance(track_name.as_str(), instance_id)
3159 .await
3160 {
3161 return;
3162 }
3163 }
3164 #[cfg(all(unix, not(target_os = "macos")))]
3165 Action::TrackShowLv2Gui { .. } => {
3166 if Self::box_bool(self.handle_track_show_lv2_gui(a.clone())).await {
3167 return;
3168 }
3169 }
3170 Action::TrackSetPluginResourceDir { .. } => {
3171 if Self::box_bool(self.handle_track_set_plugin_resource_dir(a.clone())).await {
3172 return;
3173 }
3174 }
3175 Action::TrackClapFileReferences { .. } => {
3176 if Self::box_bool(self.handle_track_clap_file_references(a.clone())).await {
3177 return;
3178 }
3179 }
3180 Action::TrackUpdateClapFileReference { .. } => {
3181 if self
3182 .handle_track_update_clap_file_reference(a.clone())
3183 .await
3184 {
3185 return;
3186 }
3187 }
3188 Action::ClipSetPluginResourceDir { .. } => {
3189 if Self::box_bool(self.handle_clip_set_plugin_resource_dir(a.clone())).await {
3190 return;
3191 }
3192 }
3193 Action::ClipClapFileReferences { .. } => {
3194 if Self::box_bool(self.handle_clip_clap_file_references(a.clone())).await {
3195 return;
3196 }
3197 }
3198 Action::ClipUpdateClapFileReference { .. } => {
3199 if Self::box_bool(self.handle_clip_update_clap_file_reference(a.clone())).await {
3200 return;
3201 }
3202 }
3203 Action::TrackSetClapParameter { .. } => {
3204 if Self::box_bool(self.handle_track_set_clap_parameter(a.clone())).await {
3205 return;
3206 }
3207 }
3208 Action::ClipSetClapParameter { .. } => {
3209 if Self::box_bool(self.handle_clip_set_clap_parameter(a.clone())).await {
3210 return;
3211 }
3212 }
3213 Action::TrackSetClapParameterAt { .. } => {
3214 if Self::box_bool(self.handle_track_set_clap_parameter_at(a.clone())).await {
3215 return;
3216 }
3217 }
3218 Action::TrackBeginClapParameterEdit { .. } => {
3219 if Self::box_bool(self.handle_track_begin_clap_parameter_edit(a.clone())).await {
3220 return;
3221 }
3222 }
3223 Action::TrackEndClapParameterEdit { .. } => {
3224 if Self::box_bool(self.handle_track_end_clap_parameter_edit(a.clone())).await {
3225 return;
3226 }
3227 }
3228 Action::TrackGetClapParameters { .. } => {
3229 if Self::box_bool(self.handle_track_get_clap_parameters(a.clone())).await {
3230 return;
3231 }
3232 }
3233 Action::TrackClapParameters { .. } => {}
3234 Action::TrackClapSnapshotState { .. } => {
3235 if Self::box_bool(self.handle_track_clap_snapshot_state(a.clone())).await {
3236 return;
3237 }
3238 }
3239 Action::ClipClapSnapshotState { .. } => {
3240 if Self::box_bool(self.handle_clip_clap_snapshot_state(a.clone())).await {
3241 return;
3242 }
3243 }
3244 Action::TrackClapStateSnapshot { .. } => {}
3245 Action::ClipClapStateSnapshot { .. } => {}
3246 Action::TrackClapStateDirty { .. } => {}
3247 Action::ClipClapStateDirty { .. } => {}
3248 Action::TrackClapRestoreState { .. } => {
3249 if Self::box_bool(self.handle_track_clap_restore_state(a.clone())).await {
3250 return;
3251 }
3252 }
3253 Action::ClipClapRestoreState { .. } => {
3254 if Self::box_bool(self.handle_clip_clap_restore_state(a.clone())).await {
3255 return;
3256 }
3257 }
3258 Action::TrackSnapshotAllClapStates { .. } => {
3259 if Self::box_bool(self.handle_track_snapshot_all_clap_states(a.clone())).await {
3260 return;
3261 }
3262 }
3263 Action::TrackSnapshotAllClapStatesDone { .. } => {}
3264 Action::TrackGetVst3Graph { .. } => {
3265 if Self::box_bool(self.handle_track_get_vst3_graph(a.clone())).await {
3266 return;
3267 }
3268 }
3269 Action::TrackVst3Graph { .. } => {}
3270 Action::TrackSetVst3Parameter { .. } => {
3271 if Self::box_bool(self.handle_track_set_vst3_parameter(a.clone())).await {
3272 return;
3273 }
3274 }
3275 Action::TrackSetPluginBypassed { .. } => {
3276 if Self::box_bool(self.handle_track_set_plugin_bypassed(a.clone())).await {
3277 return;
3278 }
3279 }
3280 Action::TrackGetVst3Parameters { .. } => {
3281 if Self::box_bool(self.handle_track_get_vst3_parameters(a.clone())).await {
3282 return;
3283 }
3284 }
3285 Action::TrackVst3Parameters { .. } => {}
3286 #[cfg(all(unix, not(target_os = "macos")))]
3287 Action::TrackSetLv2ControlValue { .. } => {
3288 if Self::box_bool(self.handle_track_set_lv2_control_value(a.clone())).await {
3289 return;
3290 }
3291 }
3292 #[cfg(all(unix, not(target_os = "macos")))]
3293 Action::TrackGetLv2PluginControls { .. } => {
3294 if Self::box_bool(self.handle_track_get_lv2_plugin_controls(a.clone())).await {
3295 return;
3296 }
3297 }
3298 #[cfg(all(unix, not(target_os = "macos")))]
3299 Action::TrackLv2SnapshotState { .. } => {
3300 if Self::box_bool(self.handle_track_lv2_snapshot_state(a.clone())).await {
3301 return;
3302 }
3303 }
3304 #[cfg(all(unix, not(target_os = "macos")))]
3305 Action::ClipLv2SnapshotState { .. } => {
3306 if Self::box_bool(self.handle_clip_lv2_snapshot_state(a.clone())).await {
3307 return;
3308 }
3309 }
3310 Action::TrackVst3SnapshotState { .. } => {
3311 if Self::box_bool(self.handle_track_vst3_snapshot_state(a.clone())).await {
3312 return;
3313 }
3314 }
3315 Action::ClipVst3SnapshotState { .. } => {
3316 if Self::box_bool(self.handle_clip_vst3_snapshot_state(a.clone())).await {
3317 return;
3318 }
3319 }
3320 Action::TrackVst3StateSnapshot { .. } => {}
3321 Action::ClipVst3StateSnapshot { .. } => {}
3322 Action::TrackVst3RestoreState { .. } => {
3323 if Self::box_bool(self.handle_track_vst3_restore_state(a.clone())).await {
3324 return;
3325 }
3326 }
3327 Action::TrackConnectVst3Audio { .. } => {
3328 if Self::box_bool(self.handle_track_connect_vst3_audio(a.clone())).await {
3329 return;
3330 }
3331 }
3332 Action::TrackDisconnectVst3Audio { .. } => {
3333 if Self::box_bool(self.handle_track_disconnect_vst3_audio(a.clone())).await {
3334 return;
3335 }
3336 }
3337 Action::ClipMove { .. } => {
3338 self.handle_clip_move(a.clone()).await;
3339 }
3340 Action::AddClip { .. } => {
3341 if Self::box_bool(self.handle_add_clip(a.clone())).await {
3342 return;
3343 }
3344 }
3345 Action::AddGroupedClip { .. } => {
3346 if Self::box_bool(self.handle_add_grouped_clip(a.clone())).await {
3347 return;
3348 }
3349 }
3350 Action::RemoveClip {
3351 ref track_name,
3352 kind,
3353 ref clip_indices,
3354 } => {
3355 self.remove_clips_from_track(track_name, kind, clip_indices);
3356 }
3357 Action::MoveClipToUnused {
3358 ref track_name,
3359 kind,
3360 ref clip_indices,
3361 } => {
3362 self.move_clips_to_unused(track_name, kind, clip_indices);
3363 }
3364 Action::DeleteUnusedClips { ref clip_ids } => {
3365 self.delete_unused_clips(clip_ids);
3366 }
3367 Action::SetUnusedClips {
3368 ref audio,
3369 ref midi,
3370 } => {
3371 self.set_unused_clips(audio.clone(), midi.clone());
3372 }
3373 Action::RenameClip {
3374 ref track_name,
3375 kind,
3376 clip_index,
3377 ref new_name,
3378 } => {
3379 self.rename_clip_references(track_name, kind, clip_index, new_name);
3380 }
3381 Action::SetClipIdentity {
3382 ref track_name,
3383 kind,
3384 clip_index,
3385 ref new_id,
3386 ref new_name,
3387 } => {
3388 self.set_clip_identity(track_name, kind, clip_index, new_id, new_name);
3389 }
3390 Action::SetClipSourceName {
3391 ref track_name,
3392 kind,
3393 clip_index,
3394 ref name,
3395 } => {
3396 self.set_clip_source_name(track_name, clip_index, kind, name.clone());
3397 }
3398 Action::SetClipFade { .. } => {
3399 if Self::box_bool(self.handle_set_clip_fade(a.clone())).await {
3400 return;
3401 }
3402 }
3403 Action::SetClipBounds {
3404 ref track_name,
3405 clip_index,
3406 kind,
3407 start,
3408 length,
3409 offset,
3410 } => {
3411 self.set_clip_bounds(track_name, clip_index, kind, start, length, offset);
3412 }
3413 Action::SyncClipBounds {
3414 ref track_name,
3415 clip_index,
3416 kind,
3417 start,
3418 length,
3419 offset,
3420 } => {
3421 self.set_clip_bounds(track_name, clip_index, kind, start, length, offset);
3422 }
3423 Action::SetClipMuted {
3424 ref track_name,
3425 clip_index,
3426 kind,
3427 muted,
3428 } => {
3429 self.set_clip_muted(track_name, clip_index, kind, muted);
3430 }
3431 Action::SetClipPluginGraphJson {
3432 ref track_name,
3433 clip_index,
3434 ref plugin_graph_json,
3435 } => {
3436 self.set_clip_plugin_graph_json(track_name, clip_index, plugin_graph_json.clone());
3437 }
3438 Action::SetClipPitchCorrection { .. } => {
3439 if Self::box_bool(self.handle_set_clip_pitch_correction(a.clone())).await {
3440 return;
3441 }
3442 }
3443 Action::Connect {
3444 ref from_track,
3445 from_port,
3446 ref to_track,
3447 to_port,
3448 kind,
3449 } => {
3450 self.handle_connect(
3451 from_track.as_str(),
3452 from_port,
3453 to_track.as_str(),
3454 to_port,
3455 kind,
3456 )
3457 .await;
3458 }
3459 Action::Disconnect { .. } => {
3460 self.handle_disconnect(a.clone()).await;
3461 }
3462 Action::OpenAudioDevice { .. } => {
3463 let (done, updated) = self.handle_open_audio_device(a.clone()).await;
3464 if done {
3465 return;
3466 }
3467 if let Some(action) = updated {
3468 action_to_process = action;
3469 }
3470 }
3471 Action::JackAddAudioInputPort => {
3472 if Self::box_bool(self.handle_jack_add_audio_input_port(a.clone())).await {
3473 return;
3474 }
3475 }
3476 Action::JackRemoveAudioInputPort(_removed_port) => {
3477 if self
3478 .handle_jack_remove_audio_input_port(_removed_port, a.clone())
3479 .await
3480 {
3481 return;
3482 }
3483 }
3484 Action::JackAddAudioOutputPort => {
3485 if Self::box_bool(self.handle_jack_add_audio_output_port(a.clone())).await {
3486 return;
3487 }
3488 }
3489 Action::JackRemoveAudioOutputPort(_removed_port) => {
3490 if self
3491 .handle_jack_remove_audio_output_port(_removed_port, a.clone())
3492 .await
3493 {
3494 return;
3495 }
3496 }
3497 Action::JackGetGraph => {
3498 #[cfg(unix)]
3499 {
3500 match self
3501 .jack_runtime
3502 .as_ref()
3503 .ok_or(
3504 "JACK runtime is not active; open the JACK backend first".to_string(),
3505 )
3506 .and_then(|jack| jack.graph_info())
3507 {
3508 Ok(graph) => self.notify_clients(Ok(Action::JackGraph(graph))).await,
3509 Err(e) => self.notify_clients(Err(e)).await,
3510 }
3511 }
3512 #[cfg(not(unix))]
3513 {
3514 self.notify_clients(Err(
3515 "JACK backend is not available on this platform build".to_string(),
3516 ))
3517 .await;
3518 }
3519 return;
3520 }
3521 Action::JackConnect {
3522 ref source,
3523 ref destination,
3524 } => {
3525 #[cfg(unix)]
3526 {
3527 match self
3528 .jack_runtime
3529 .as_ref()
3530 .ok_or(
3531 "JACK runtime is not active; open the JACK backend first".to_string(),
3532 )
3533 .and_then(|jack| jack.connect_ports_by_name(source, destination))
3534 .and_then(|_| {
3535 self.jack_runtime
3536 .as_ref()
3537 .expect("JACK runtime was checked")
3538 .graph_info()
3539 }) {
3540 Ok(graph) => {
3541 self.notify_clients(Ok(a.clone())).await;
3542 self.notify_clients(Ok(Action::JackGraph(graph))).await;
3543 }
3544 Err(e) => self.notify_clients(Err(e)).await,
3545 }
3546 }
3547 #[cfg(not(unix))]
3548 {
3549 self.notify_clients(Err(
3550 "JACK backend is not available on this platform build".to_string(),
3551 ))
3552 .await;
3553 }
3554 return;
3555 }
3556 Action::JackDisconnect {
3557 ref source,
3558 ref destination,
3559 } => {
3560 #[cfg(unix)]
3561 {
3562 match self
3563 .jack_runtime
3564 .as_ref()
3565 .ok_or(
3566 "JACK runtime is not active; open the JACK backend first".to_string(),
3567 )
3568 .and_then(|jack| jack.disconnect_ports_by_name(source, destination))
3569 .and_then(|_| {
3570 self.jack_runtime
3571 .as_ref()
3572 .expect("JACK runtime was checked")
3573 .graph_info()
3574 }) {
3575 Ok(graph) => {
3576 self.notify_clients(Ok(a.clone())).await;
3577 self.notify_clients(Ok(Action::JackGraph(graph))).await;
3578 }
3579 Err(e) => self.notify_clients(Err(e)).await,
3580 }
3581 }
3582 #[cfg(not(unix))]
3583 {
3584 self.notify_clients(Err(
3585 "JACK backend is not available on this platform build".to_string(),
3586 ))
3587 .await;
3588 }
3589 return;
3590 }
3591 Action::OpenMidiInputDevice(ref device) => {
3592 if let Some(worker) = &self.hw_worker {
3593 if let Err(e) = worker
3594 .tx
3595 .send(Message::HWOpenMidiInputDevice(device.clone()))
3596 .await
3597 {
3598 self.notify_clients(Err(format!("Failed to send MIDI input open: {e}")))
3599 .await;
3600 }
3601 return;
3602 }
3603 let Some(midi_hub) = self.midi_hub.as_mut() else {
3604 self.notify_clients(Err("Hardware MIDI hub is not available".to_string()))
3605 .await;
3606 return;
3607 };
3608 if let Err(e) = midi_hub.open_input(device) {
3609 self.notify_clients(Err(e)).await;
3610 return;
3611 }
3612 }
3613 Action::OpenMidiOutputDevice(ref device) => {
3614 if let Some(worker) = &self.hw_worker {
3615 if let Err(e) = worker
3616 .tx
3617 .send(Message::HWOpenMidiOutputDevice(device.clone()))
3618 .await
3619 {
3620 self.notify_clients(Err(format!("Failed to send MIDI output open: {e}")))
3621 .await;
3622 }
3623 return;
3624 }
3625 let Some(midi_hub) = self.midi_hub.as_mut() else {
3626 self.notify_clients(Err("Hardware MIDI hub is not available".to_string()))
3627 .await;
3628 return;
3629 };
3630 if let Err(e) = midi_hub.open_output(device) {
3631 self.notify_clients(Err(e)).await;
3632 return;
3633 }
3634 }
3635 Action::RequestSessionDiagnostics => {
3636 self.handle_request_session_diagnostics().await;
3637 }
3638 Action::RequestMidiLearnMappingsReport => {
3639 self.handle_request_midi_learn_mappings_report().await;
3640 }
3641 Action::ClearAllMidiLearnBindings => {
3642 if Self::box_bool(self.handle_clear_all_midi_learn_bindings(a.clone())).await {
3643 return;
3644 }
3645 }
3646 #[cfg(all(unix, not(target_os = "macos")))]
3647 Action::TrackLv2PluginControls { .. } => {}
3648 #[cfg(all(unix, not(target_os = "macos")))]
3649 Action::ClipLv2PluginControls { .. } => {}
3650 #[cfg(all(unix, not(target_os = "macos")))]
3651 Action::TrackLv2StateSnapshot { .. } => {}
3652 #[cfg(all(unix, not(target_os = "macos")))]
3653 Action::ClipLv2StateSnapshot { .. } => {}
3654 #[cfg(all(unix, not(target_os = "macos")))]
3655 Action::TrackLv2Midnam { .. } => {}
3656 Action::TrackClapNoteNames { .. } => {}
3657 Action::SessionDiagnosticsReport { .. } => {}
3658 Action::MidiLearnMappingsReport { .. } => {}
3659 Action::HWInfo { .. } => {}
3660 Action::HistoryState { .. } => {}
3661 Action::Undo => {}
3662 Action::Redo => {}
3663 Action::ApplyGroupedActions(_) => {}
3664 _ => {}
3665 }
3666
3667 if let Some(inverse) = inverse_actions {
3668 if let Some(group) = self.history_group.as_mut() {
3669 group.forward_actions.push(action_to_process.clone());
3670 group.inverse_actions.splice(0..0, inverse);
3671 } else {
3672 self.history.record(UndoEntry {
3673 forward_actions: vec![action_to_process.clone()],
3674 inverse_actions: inverse,
3675 });
3676 }
3677 }
3678
3679 self.notify_clients(Ok(action_to_process)).await;
3680 }
3681 pub async fn work(&mut self) {
3682 let mut tick = tokio::time::interval(Duration::from_millis(1));
3685 tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
3686 loop {
3687 let message = tokio::select! {
3688 message = self.rx.recv() => {
3689 let Some(message) = message else {
3690 break;
3691 };
3692 message
3693 }
3694 _ = tick.tick() => {
3695 self.poll_node_worker_results().await;
3696 self.poll_jack_hw_finished().await;
3697 self.poll_stopped_plugin_parameter_echoes().await;
3698 self.on_executor_tick().await;
3699 continue;
3700 }
3701 };
3702 self.poll_node_worker_results().await;
3703 self.poll_jack_hw_finished().await;
3704 self.poll_stopped_plugin_parameter_echoes().await;
3705 match message {
3706 Message::Ready(id) => {
3707 if let Some(track_name) = self.bounce_worker_tracks.remove(&id) {
3711 self.offline_bounce_jobs.remove(&track_name);
3712 }
3713 self.push_ready_worker(id);
3714 if self.dispatch_node_jobs(Vec::new()).await {
3715 self.on_all_tracks_finished().await;
3716 }
3717 self.drain_pending_requests_if_idle().await;
3718 }
3719 Message::NodeDone {
3720 worker_id,
3721 epoch,
3722 node,
3723 output_linear,
3724 parameter_updates,
3725 latency_changed,
3726 } => {
3727 self.on_node_done(
3728 worker_id,
3729 epoch,
3730 node,
3731 output_linear,
3732 parameter_updates,
3733 latency_changed,
3734 )
3735 .await;
3736 }
3737 Message::Channel(s) => {
3738 self.clients.push(s);
3739 }
3740 Message::Response(result) => {
3741 self.notify_clients(result).await;
3742 }
3743
3744 Message::Request(a) => {
3745 self.dispatch_request(a).await;
3746 self.plan_builder.mark_dirty();
3749 }
3750 Message::OscRequest { action, reply_to } => {
3751 self.osc_reply_target = Some(reply_to);
3752 self.dispatch_request(action).await;
3753 self.osc_reply_target = None;
3754 self.plan_builder.mark_dirty();
3755 }
3756 Message::OfflineBounceFinished { result } => {
3757 if let Ok(Action::TrackOfflineBounce { track_name, .. })
3758 | Ok(Action::TrackOfflineBounceCanceled { track_name, .. }) = &result
3759 {
3760 self.offline_bounce_jobs.remove(track_name);
3761 }
3762 self.notify_clients(result).await;
3763 self.drain_pending_requests_if_idle().await;
3764 }
3765 Message::HWFinished => {
3766 if !self.awaiting_hwfinished {
3767 tracing::debug!("HWFinished ignored (not awaiting)");
3768 continue;
3769 }
3770 tracing::debug!("HWFinished handling; playing={}", self.playing);
3771 self.handling_hwfinished = true;
3772 self.awaiting_hwfinished = false;
3773 #[cfg(unix)]
3774 {
3775 if let Some(jack) = self.jack_runtime.as_mut() {
3776 if !self.pending_hw_midi_out_events.is_empty() {
3777 let out_events =
3778 std::mem::take(&mut self.pending_hw_midi_out_events);
3779 jack.write_events(&out_events);
3780 }
3781 let mut in_events = vec![];
3782 jack.read_events_into(&mut in_events);
3783 if !in_events.is_empty() {
3784 self.pending_hw_midi_events.extend(in_events);
3785 }
3786 let dropped = jack.take_midi_events_dropped();
3787 if dropped > 0 {
3788 tracing::warn!(
3789 "JACK MIDI ring full; {dropped} events dropped since last cycle"
3790 );
3791 }
3792 }
3793 }
3794 #[cfg(unix)]
3795 if self.jack_runtime.is_some() {
3796 self.sync_from_jack_transport().await;
3797 }
3798 while let Some(a) = self.pending_requests.pop_front() {
3799 self.handle_request(a).await;
3800 }
3801 self.apply_mute_solo_policy();
3802 self.append_recorded_cycle();
3803 self.flush_completed_recordings().await;
3804 let hw_in_routes = self.midi_hw_in_routes.clone();
3805 let pending_hw_in_by_device = self.pending_hw_midi_events_by_device.clone();
3806 let mut reconfigured_tracks = Vec::new();
3807 let state = self.state_snapshot.load_full();
3808 for (track_name, track) in state.tracks.iter() {
3809 let mut track_lock = track.lock();
3810 if self.jack_runtime_is_some() {
3811 if !self.pending_hw_midi_events.is_empty() {
3812 track_lock.push_hw_midi_events(&self.pending_hw_midi_events);
3813 }
3814 } else {
3815 for route in hw_in_routes.iter().filter(|r| &r.to_track == track_name) {
3816 if let Some(events) = pending_hw_in_by_device.get(&route.device) {
3817 track_lock.push_hw_midi_events_to_port(route.to_port, events);
3818 }
3819 }
3820 }
3821 if track_lock.setup() {
3822 reconfigured_tracks.push(track_name.clone());
3823 }
3824 }
3825 self.publish_track_meters();
3826 self.publish_session_runtime_reports().await;
3827 self.publish_clap_state_dirty().await;
3828 for track_name in reconfigured_tracks {
3829 let track = state.tracks.get(&track_name).cloned();
3830 if let Some(track) = track {
3831 let (plugins, connections, connectable_connections) = {
3832 let track_lock = track.lock();
3833 (
3834 track_lock.plugin_graph_plugins(false),
3835 track_lock.plugin_graph_connections(),
3836 track_lock.connectable_connections(),
3837 )
3838 };
3839 self.notify_clients(Ok(Action::TrackPluginGraph {
3840 track_name: track_name.clone(),
3841 plugins,
3842 connections,
3843 connectable_connections,
3844 }))
3845 .await;
3846 }
3847 }
3848 self.pending_hw_midi_events.clear();
3849 self.pending_hw_midi_events_by_device.clear();
3850 if self.transport_running {
3851 if self.transport_panic_flush_pending {
3852 self.transport_panic_flush_pending = false;
3853 } else if self.transport_restart_pending {
3854 self.transport_restart_pending = false;
3855 } else {
3856 let next = self
3857 .transport_sample
3858 .saturating_add(self.current_cycle_samples());
3859 let normalized = self.normalize_transport_sample(next);
3860 let wrapped = normalized != next;
3861 self.transport_sample = normalized;
3862 self.publish_transport_snapshot();
3863 if wrapped {
3864 if self.notified_loop_wrap_sample == Some(self.transport_sample) {
3865 self.notified_loop_wrap_sample = None;
3866 } else {
3867 self.notify_clients(Ok(Action::TransportPosition(
3868 self.transport_sample,
3869 )))
3870 .await;
3871 }
3872 }
3873 }
3874 }
3875 if self.session_clip_playback_enabled && self.playing {
3876 self.session_transport_sample = self
3877 .session_transport_sample
3878 .saturating_add(self.current_cycle_samples());
3879 }
3880 {
3881 let echoes = self.apply_modulators(self.active_transport_sample());
3882 for action in echoes {
3883 self.notify_clients(Ok(action)).await;
3884 }
3885 }
3886 self.start_plan_cycle().await;
3887 #[cfg(unix)]
3888 {
3889 if self.jack_runtime.is_some() {
3890 self.awaiting_hwfinished = true;
3891 }
3892 }
3893 self.handling_hwfinished = false;
3894 }
3895 Message::HWMidiEvents(events) => {
3896 for hw_event in events {
3897 let thru_targets: Vec<String> = self
3898 .midi_hw_thru_routes
3899 .iter()
3900 .filter(|route| route.from_device == hw_event.device)
3901 .map(|route| route.to_device.clone())
3902 .collect();
3903 for device in thru_targets {
3904 self.pending_hw_midi_out_events_by_device.push(HwMidiEvent {
3905 device,
3906 event: hw_event.event.clone(),
3907 });
3908 }
3909 if hw_event.event.data.len() >= 3 {
3910 let status = hw_event.event.data[0];
3911 if status & 0xF0 == 0xB0 {
3912 let channel = status & 0x0F;
3913 let cc = hw_event.event.data[1];
3914 let value = hw_event.event.data[2];
3915 self.handle_incoming_hw_cc(&hw_event.device, channel, cc, value)
3916 .await;
3917 }
3918 if self.step_recording_enabled && status & 0xF0 == 0x90 {
3919 let channel = status & 0x0F;
3920 let pitch = hw_event.event.data[1];
3921 let velocity = hw_event.event.data[2];
3922 if velocity > 0 {
3923 self.notify_clients(Ok(Action::StepRecordMidiNote {
3924 device: hw_event.device.clone(),
3925 channel,
3926 pitch,
3927 velocity,
3928 }))
3929 .await;
3930 }
3931 }
3932 }
3933 self.pending_hw_midi_events_by_device
3934 .entry(hw_event.device)
3935 .or_default()
3936 .push(hw_event.event);
3937 }
3938 }
3939 _ => {}
3940 }
3941 }
3942 }
3943
3944 pub(crate) fn collect_hw_midi_output_events(&self) -> Vec<MidiEvent> {
3945 let mut events = vec![];
3946 for track in self.state_snapshot.load_full().tracks.values() {
3947 events.extend(
3948 track
3949 .lock()
3950 .take_hw_midi_out_events()
3951 .into_iter()
3952 .map(|evt| evt.event),
3953 );
3954 }
3955 events.sort_by_key(|a| a.frame);
3956 events
3957 }
3958
3959 pub(crate) fn collect_hw_midi_output_events_by_device(&mut self) -> Vec<HwMidiEvent> {
3960 let mut events = Vec::<HwMidiEvent>::new();
3961 let routes = self.midi_hw_out_routes.clone();
3962 let mut events_by_track = HashMap::<String, Vec<crate::track::HwMidiOutEvent>>::new();
3963 {
3964 let state = self.state_snapshot.load_full();
3965 for route in &routes {
3966 if events_by_track.contains_key(&route.from_track) {
3967 continue;
3968 }
3969 let Some(track) = state.tracks.get(&route.from_track) else {
3970 continue;
3971 };
3972 events_by_track.insert(
3973 route.from_track.clone(),
3974 track.lock().take_hw_midi_out_events(),
3975 );
3976 }
3977 }
3978
3979 for route in routes {
3980 let Some(track_events) = events_by_track.get(&route.from_track) else {
3981 continue;
3982 };
3983 for hw_event in track_events
3984 .iter()
3985 .filter(|evt| evt.port == route.from_port)
3986 {
3987 self.update_active_hw_notes_for_track(
3988 &route.from_track,
3989 &route.device,
3990 &hw_event.event.data,
3991 );
3992 events.push(HwMidiEvent {
3993 device: route.device.clone(),
3994 event: hw_event.event.clone(),
3995 });
3996 }
3997 }
3998 events.sort_by(|a, b| {
3999 a.event
4000 .frame
4001 .cmp(&b.event.frame)
4002 .then_with(|| a.device.cmp(&b.device))
4003 });
4004 events
4005 }
4006}